diff --git a/netx_api/biz_state/claim.py b/netx_api/biz_state/claim.py index 7eaab36..a86674e 100644 --- a/netx_api/biz_state/claim.py +++ b/netx_api/biz_state/claim.py @@ -39,6 +39,13 @@ def enqueue_collect(task_id: str, *, manual: bool = False) -> dict[str, Any]: task = db.get(BizStateTask, tid) if not task: return {"ok": False, "queued": False, "reason": "task_not_found", "task_id": tid} + if str(task.source or "").strip().lower() == "import": + return { + "ok": False, + "queued": False, + "reason": "import_offline_only", + "task_id": tid, + } if bool(task.collect_running): return { "ok": True, diff --git a/netx_api/biz_state/import_runner.py b/netx_api/biz_state/import_runner.py new file mode 100644 index 0000000..c31e6ef --- /dev/null +++ b/netx_api/biz_state/import_runner.py @@ -0,0 +1,662 @@ +"""Offline log import → same biz_state batch / parse / persist path (no SSH).""" + +from __future__ import annotations + +import logging +import threading +from pathlib import Path +from typing import Any +from uuid import uuid4 + +from fastapi import HTTPException + +from ..db import SessionLocal +from ..lldp_shared import resolve_vendor_key +from ..models import BizStateBatch, BizStateTask +from ..timeutil import utcnow_naive +from .collect_runner import ( + _finish_task, + _finalize_batch_status, + _format_error, + _run_primary_parse_job, +) +from .collect_session import resolve_aux_command +from .command_match import match_command, normalize_command +from .log_split import LogSegment, unpack_upload +from .parse_pool import AuxRawCapture, PrimaryParseJob +from .parsers import get_parser +from .persist_pool import get_persist_pool +from .spool import ( + SpooledCommand, + clear_batch_spool, + count_text_lines, + persist_every_cmds, + write_meta, + write_raw_text, +) + +_log = logging.getLogger("netx.biz_state.import") + + +def _utcnow(): + return utcnow_naive() + + +def _vendor_fields_for_key(vendor_key: str) -> tuple[str, str]: + """Map vendor_key → (vendor label, device_type) for offline import tasks.""" + key = str(vendor_key or "").strip().lower() or "zte" + if key.startswith("huawei") or key in ("vrp", "ce", "ne"): + return "Huawei", "huawei_vrp" + if key.startswith("cisco") or key in ("ios", "nxos", "iosxe", "iosxr"): + return "Cisco", "cisco_ios" + if key.startswith("zte") or key in ("zxros", "zxr10"): + return "ZTE", "zte_zxros" + if key in ("generic", "any", "*"): + return "generic", "" + return key[:64] or "generic", "" + + +def create_standalone_import_task( + *, + vendor_key: str, + ne_name: str = "", + note: str = "", + filename: str = "", +) -> dict[str, Any]: + """Create a paused offline task (source=import) with no inventory NE.""" + vk = str(vendor_key or "").strip().lower() or "zte" + vendor, device_type = _vendor_fields_for_key(vk) + fname = Path(str(filename or "")).name[:120] + label = str(ne_name or "").strip() or (fname and f"import:{fname}") or "offline-import" + tid = uuid4().hex + ne_id = f"import-{tid[:12]}" + now = _utcnow() + db = SessionLocal() + try: + task = BizStateTask( + id=tid, + source="import", + ne_id=ne_id, + ne_name=label[:256], + ne_ip="", + vendor=vendor, + device_type=device_type, + note=str(note or "")[:256], + purpose="", + status="paused", # offline only — no schedule / SSH + interval_sec=3600, + retention_days=30, + daily_keep_enabled=False, + daily_keep_count=10, + retention_batches=30, + created_at=now, + updated_at=now, + ) + db.add(task) + db.commit() + return { + "ok": True, + "task_id": tid, + "ne_id": ne_id, + "ne_name": task.ne_name, + "vendor": vendor, + "device_type": device_type, + "vendor_key": vk, + } + except Exception: + _log.exception("create_standalone_import_task failed") + try: + db.rollback() + except Exception: + pass + raise + finally: + db.close() + + +def enqueue_import( + task_id: str, + *, + filename: str = "", + vendor_key: str = "", +) -> dict[str, Any]: + """Create a running import batch; mark task collect_running (mutex with SSH collect).""" + tid = str(task_id or "").strip() + if not tid: + return {"ok": False, "queued": False, "reason": "missing_task_id", "task_id": ""} + + db = SessionLocal() + try: + task = db.get(BizStateTask, tid) + if not task: + return {"ok": False, "queued": False, "reason": "task_not_found", "task_id": tid} + if bool(task.collect_running): + return { + "ok": True, + "queued": False, + "reason": "already_collecting", + "task_id": tid, + } + st = str(task.status or "").strip() + if st in ("", "deleted"): + return {"ok": False, "queued": False, "reason": "bad_status", "task_id": tid} + + now = _utcnow() + task.collect_running = True + task.last_collect_started_at = now + task.last_error = "" + task.updated_at = now + if hasattr(task, "collect_queued_at"): + task.collect_queued_at = now + + fname = Path(str(filename or "import.log")).name[:120] or "import.log" + vk = str(vendor_key or "").strip().lower() + if not vk: + vk = resolve_vendor_key(task.vendor or "", task.device_type or "") + + batch = BizStateBatch( + id=uuid4().hex, + task_id=task.id, + source=task.source, + ne_id=task.ne_id, + ne_name=task.ne_name, + vendor=task.vendor, + status="running", + started_at=now, + message=f"importing:{fname}", + alias=fname[:128], + ) + db.add(batch) + db.commit() + return { + "ok": True, + "queued": True, + "batch_id": batch.id, + "task_id": tid, + "vendor_key": vk, + "filename": fname, + "manual": True, + "import": True, + } + except Exception: + _log.exception("enqueue_import failed task=%s", tid) + try: + db.rollback() + except Exception: + pass + return {"ok": False, "queued": False, "reason": "enqueue_failed", "task_id": tid} + finally: + db.close() + + +def _segment_index(segments: list[LogSegment]) -> dict[str, LogSegment]: + """Map normalize_command(cmd) → first matching segment.""" + idx: dict[str, LogSegment] = {} + for seg in segments: + key = normalize_command(seg.command) + if key and key not in idx: + idx[key] = seg + return idx + + +def _build_aux_captures( + *, + profile: Any, + params: dict[str, str], + batch_id: str, + seg_index: dict[str, LogSegment], +) -> tuple[list[AuxRawCapture], list[str]]: + """Try to satisfy aux_commands from segments already in the import log.""" + captures: list[AuxRawCapture] = [] + notes: list[str] = [] + for aux in list(getattr(profile, "aux_commands", None) or []): + try: + ra = resolve_aux_command(aux, params) + except Exception as exc: + notes.append(f"aux:{getattr(aux, 'key', '?')}:resolve:{exc}") + continue + hit = seg_index.get(normalize_command(ra.command)) + if hit is None: + notes.append(f"enrich_skipped:missing_aux:{ra.key}") + continue + aux_id = uuid4().hex + raw_rel = "" + try: + raw_rel = write_raw_text(batch_id, aux_id, hit.body) + except Exception: + _log.exception("import aux spool raw failed") + captures.append( + AuxRawCapture( + key=ra.key, + aux_id=aux_id, + profile_id=ra.profile_id, + parser_id=ra.parser_id, + metric_id=str(getattr(ra.profile, "metric_id", "") or ""), + command=ra.command, + textfsm_command=ra.textfsm_command or ra.command, + rule_keys=tuple(ra.rule_keys or ()), + raw=hit.body, + raw_rel_path=raw_rel, + cache_hit=False, + ok=True, + ) + ) + return captures, notes + + +def run_import_batch( + *, + batch_id: str, + task_id: str, + segments: list[LogSegment], + vendor_key: str, + vendor: str = "", + device_type: str = "", + filename: str = "", +) -> dict[str, Any]: + """Parse imported segments into the batch (sync parse + persist pool).""" + bid = str(batch_id or "").strip() + tid = str(task_id or "").strip() + vk = str(vendor_key or "").strip().lower() or "zte" + any_ok = False + any_fail = False + matched = 0 + unmatched = 0 + cmd_count = 0 + lane_errors: list[str] = [] + + try: + clear_batch_spool(bid) + except Exception: + _log.exception("import clear spool failed batch=%s", bid) + + persist = get_persist_pool() + pending: list[SpooledCommand] = [] + flush_every = persist_every_cmds() + seg_index = _segment_index(segments) + persisted: set[tuple[str, str]] = set() + cache_lock = threading.RLock() + + def _submit_pending() -> None: + nonlocal pending + if not pending: + return + chunk = list(pending) + pending = [] + persist.submit(bid, chunk) + + def _queue(item: SpooledCommand) -> None: + nonlocal cmd_count + try: + write_meta(bid, item.id, item.to_meta()) + except Exception: + _log.exception("import write meta failed cmd=%s", item.id) + pending.append(item) + cmd_count += 1 + if len(pending) >= flush_every: + _submit_pending() + + for seg in segments: + cmd = normalize_command(seg.command) + cmd_id = uuid4().hex + src_note = "" + if seg.source_file: + src_note = f"file={seg.source_file}" + hit = match_command(vendor_key=vk, command=cmd) + raw_rel = "" + try: + raw_rel = write_raw_text(bid, cmd_id, seg.body) + except Exception: + _log.exception("import spool raw failed cmd=%s", cmd_id) + + if not hit: + unmatched += 1 + any_fail = True + msg = "no profile matched concrete command" + if src_note: + msg = f"{msg};{src_note}" + _queue( + SpooledCommand( + id=cmd_id, + batch_id=bid, + profile_id="", + raw_command=cmd[:512], + parse_status="unmatched", + message=msg[:1020], + raw_rel_path=raw_rel, + raw_line_count=count_text_lines(seg.body), + ) + ) + continue + + if not get_parser(hit.profile.parser_id): + unmatched += 1 + any_fail = True + _queue( + SpooledCommand( + id=cmd_id, + batch_id=bid, + profile_id=hit.profile.profile_id, + parser_id=hit.profile.parser_id, + metric_id=hit.profile.metric_id, + raw_command=cmd[:512], + params_json=dict(hit.params or {}), + parse_status="failed", + message=f"unknown parser {hit.profile.parser_id}"[:1020], + raw_rel_path=raw_rel, + raw_line_count=count_text_lines(seg.body), + ) + ) + continue + + matched += 1 + aux_caps, aux_notes = _build_aux_captures( + profile=hit.profile, + params=dict(hit.params or {}), + batch_id=bid, + seg_index=seg_index, + ) + enrich = list(hit.profile.enrich_joins or []) + # Without in-log aux, skip enrich so we don't pretend peer intent joined. + if any("enrich_skipped" in n for n in aux_notes): + enrich = [] + job = PrimaryParseJob( + batch_id=bid, + cmd_id=cmd_id, + task_item_id="", + profile_id=hit.profile.profile_id, + parser_id=hit.profile.parser_id, + metric_id=hit.profile.metric_id, + concrete=cmd, + merged_params=dict(hit.params or {}), + raw_text=seg.body, + raw_rel_path=raw_rel, + raw_line_count=count_text_lines(seg.body), + textfsm_command=hit.profile.textfsm_command or cmd, + vendor=vendor, + device_type=device_type, + enrich_joins=enrich, + aux_captures=aux_caps, + persisted=persisted, + cache_lock=cache_lock, + on_done=None, + ) + try: + ok, fail = _run_primary_parse_job(job) + if ok: + any_ok = True + if fail: + any_fail = True + if aux_notes: + _log.info( + "import enrich notes batch=%s cmd=%s %s", + bid, + cmd_id, + ";".join(aux_notes), + ) + cmd_count += 1 + len(aux_caps) + except Exception as exc: + any_fail = True + _queue( + SpooledCommand( + id=cmd_id, + batch_id=bid, + profile_id=hit.profile.profile_id, + parser_id=hit.profile.parser_id, + metric_id=hit.profile.metric_id, + raw_command=cmd[:512], + params_json=dict(hit.params or {}), + parse_status="failed", + message=f"parse: {_format_error(exc)}"[:1020], + raw_rel_path=raw_rel, + raw_line_count=count_text_lines(seg.body), + ) + ) + + _submit_pending() + # _run_primary_parse_job also submits to this same pool — wait once for all. + if not persist.wait_idle(timeout=3600.0): + any_fail = True + lane_errors.append("import: persist_barrier_timeout") + + fname = Path(str(filename or "")).name + summary = ( + f"imported:{fname or 'upload'}; segments={len(segments)}; " + f"matched={matched}; unmatched={unmatched}" + ) + if fname: + lane_errors.insert(0, summary) + + status = "" + try: + status = _finalize_batch_status( + batch_id=bid, + task_id=tid, + cmd_count=cmd_count, + total_rows=0, + any_fail=any_fail, + any_ok=any_ok or matched > 0, + lane_errors=lane_errors, + ) + except Exception as exc: + _log.exception("import finalize failed batch=%s", bid) + lane_errors.append(_format_error(exc)) + try: + from .collect_runner import _fail_batch_status + + _fail_batch_status(bid, "; ".join(lane_errors)[:1020]) + except Exception: + pass + status = "failed" + + # Full success clears message in finalize — keep import provenance visible. + if status == "success": + try: + _stamp_import_message(bid, summary) + except Exception: + _log.exception("import stamp message failed batch=%s", bid) + + err = "" + if status in ("failed",) and not any_ok: + err = "; ".join(lane_errors)[:1020] or "import failed" + try: + _finish_task(tid, error=err) + except Exception: + _log.exception("import finish task failed task=%s", tid) + + try: + clear_batch_spool(bid) + except Exception: + pass + + return { + "ok": status in ("success", "partial"), + "batch_id": bid, + "task_id": tid, + "status": status, + "segments": len(segments), + "matched": matched, + "unmatched": unmatched, + "command_count": cmd_count, + } + + +def _stamp_import_message(batch_id: str, summary: str) -> None: + db = SessionLocal() + try: + batch = db.get(BizStateBatch, batch_id) + if not batch: + return + if not str(batch.message or "").strip(): + batch.message = str(summary or "")[:1020] + db.commit() + finally: + db.close() + + +def start_standalone_import( + *, + filename: str, + data: bytes, + vendor_key: str = "", + ne_name: str = "", + note: str = "", +) -> dict[str, Any]: + """Create offline task + unpack + enqueue import (no inventory NE).""" + vk = str(vendor_key or "").strip().lower() or "zte" + try: + segments, stats = unpack_upload( + filename=filename, data=data, vendor_key=vk + ) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + if not segments: + raise HTTPException( + status_code=400, + detail="no show/display command segments found in upload", + ) + + created = create_standalone_import_task( + vendor_key=vk, + ne_name=ne_name, + note=note, + filename=filename, + ) + tid = str(created["task_id"]) + vendor = str(created.get("vendor") or "") + device_type = str(created.get("device_type") or "") + + result = enqueue_import(tid, filename=filename, vendor_key=vk) + if not result.get("queued"): + # Best-effort: leave the empty task for the user to retry / delete. + return { + "ok": bool(result.get("ok", False)), + "started": False, + "queued": False, + "reason": result.get("reason") or "enqueue_failed", + "task_id": tid, + "created_task": True, + **stats, + } + + return { + "ok": True, + "started": True, + "queued": True, + "batch_id": result["batch_id"], + "task_id": tid, + "vendor_key": vk, + "filename": result.get("filename") or filename, + "vendor": vendor, + "device_type": device_type, + "ne_name": created.get("ne_name") or "", + "created_task": True, + "segments_preview": stats, + "segments": segments, + "collect_running": True, + } + + +def start_import_from_upload( + *, + task_id: str, + filename: str, + data: bytes, + vendor_key: str = "", +) -> dict[str, Any]: + """Enqueue + unpack; caller should run ``execute_import`` in background.""" + tid = str(task_id or "").strip() + db = SessionLocal() + try: + task = db.get(BizStateTask, tid) + if not task: + raise HTTPException(status_code=404, detail="task not found") + vendor = str(task.vendor or "") + device_type = str(task.device_type or "") + vk = str(vendor_key or "").strip().lower() + if not vk: + vk = resolve_vendor_key(vendor, device_type) + finally: + db.close() + + try: + segments, stats = unpack_upload( + filename=filename, data=data, vendor_key=vk + ) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + if not segments: + raise HTTPException( + status_code=400, + detail="no show/display command segments found in upload", + ) + + result = enqueue_import(tid, filename=filename, vendor_key=vk) + if not result.get("queued"): + return { + "ok": bool(result.get("ok", False)), + "started": False, + "queued": False, + "reason": result.get("reason") or "enqueue_failed", + "task_id": tid, + **stats, + } + + return { + "ok": True, + "started": True, + "queued": True, + "batch_id": result["batch_id"], + "task_id": tid, + "vendor_key": vk, + "filename": result.get("filename") or filename, + "vendor": vendor, + "device_type": device_type, + "segments_preview": stats, + "segments": segments, # passed to background runner (in-memory) + "collect_running": True, + } + + +def execute_import( + *, + batch_id: str, + task_id: str, + segments: list[LogSegment], + vendor_key: str, + vendor: str = "", + device_type: str = "", + filename: str = "", +) -> dict[str, Any]: + """Background entry: parse imported segments.""" + try: + return run_import_batch( + batch_id=batch_id, + task_id=task_id, + segments=segments, + vendor_key=vendor_key, + vendor=vendor, + device_type=device_type, + filename=filename, + ) + except Exception as exc: + _log.exception("execute_import failed batch=%s", batch_id) + try: + from .collect_runner import _fail_batch_status + + _fail_batch_status(batch_id, _format_error(exc)) + except Exception: + pass + try: + _finish_task(task_id, error=_format_error(exc)) + except Exception: + pass + return { + "ok": False, + "batch_id": batch_id, + "task_id": task_id, + "status": "failed", + "error": _format_error(exc), + } diff --git a/netx_api/biz_state/log_split.py b/netx_api/biz_state/log_split.py new file mode 100644 index 0000000..1532d04 --- /dev/null +++ b/netx_api/biz_state/log_split.py @@ -0,0 +1,241 @@ +"""Split device CLI transcript logs into show/display command segments.""" + +from __future__ import annotations + +import io +import re +import zipfile +from dataclasses import dataclass +from pathlib import Path +from typing import BinaryIO + +from ..config import settings + +# Text-like extensions accepted from zip members / bare uploads. +_TEXT_SUFFIXES = {".txt", ".log", ".ini", ".cfg", ".cli", ".out", ".text"} + + +@dataclass(frozen=True) +class LogSegment: + """One CLI command and its captured output body.""" + + command: str + body: str + source_file: str = "" + line_start: int = 0 + + +def import_max_bytes() -> int: + return max( + 1, + int(getattr(settings, "biz_state_import_max_bytes", 256 * 1024 * 1024) or 0) + or (256 * 1024 * 1024), + ) + + +def import_max_files() -> int: + return max(1, int(getattr(settings, "biz_state_import_max_files", 200) or 200)) + + +def normalize_log_text(text: str) -> str: + """NBSP → space, unify newlines, strip trailing CR.""" + s = str(text or "").replace("\u00a0", " ").replace("\r\n", "\n").replace("\r", "\n") + return s + + +def _anchor_re(vendor_key: str) -> re.Pattern[str]: + key = str(vendor_key or "").strip().lower() + if key.startswith("huawei") or key in ("vrp", "ce", "ne"): + # Prefer display; also accept show (mixed dumps). + return re.compile(r"(?im)^(?:display|show)\s+") + if key.startswith("cisco") or key in ("ios", "nxos", "iosxe", "iosxr"): + return re.compile(r"(?im)^(?:show)\s+") + if key.startswith("zte") or key in ("zxros", "zxr10"): + return re.compile(r"(?im)^(?:show)\s+") + # Unknown / generic: both + return re.compile(r"(?im)^(?:show|display)\s+") + + +_PROMPT_LINE_RE = re.compile(r"^[A-Za-z0-9._\-\[\]/]+[#>]\s*(.+)$") +_HW_PROMPT_LINE_RE = re.compile(r"^<[^>]+>\s*(.+)$") + + +def _strip_prompt_noise(line: str) -> str: + """Drop hostname# / hostname> / prefixes from a command line.""" + s = line.strip() + if not s: + return "" + # Whole-line prompt alone + if re.fullmatch(r"[A-Za-z0-9._\-\[\]/]+[#>]", s): + return "" + if re.fullmatch(r"<[^>]+>", s): + return "" + # "R1#show arp" → "show arp" + m = _PROMPT_LINE_RE.match(s) + if m: + return m.group(1).strip() + # "display ip routing-table" + m = _HW_PROMPT_LINE_RE.match(s) + if m: + return m.group(1).strip() + return s + + +def split_log_text( + text: str, + *, + vendor_key: str = "", + source_file: str = "", +) -> list[LogSegment]: + """Split a CLI transcript into show/display segments. + + Handles real device pastes such as:: + + MDN-BCP-CN1-ZM8SP#show arp | one-line + ... + MDN-BCP-CN1-ZM8SP#show interface brief + ... + + Prompt prefixes (``host#`` / ``host>`` / ````) are stripped before + matching; each new show/display line starts a new segment. + + Anything before the first show/display (banners, clocks, lone prompts) is + discarded. If the whole text has no show/display, returns an empty list. + """ + raw = normalize_log_text(text) + if not raw.strip(): + return [] + anchor = _anchor_re(vendor_key) + lines = raw.split("\n") + starts: list[tuple[int, str]] = [] # (0-based line idx, command) + for i, line in enumerate(lines): + cleaned = _strip_prompt_noise(line) + if not cleaned: + continue + if anchor.match(cleaned): + # Command is the cleaned line (may include | filters) + cmd = re.sub(r"\s+", " ", cleaned).strip() + starts.append((i, cmd)) + + if not starts: + return [] + + out: list[LogSegment] = [] + for idx, (line_i, cmd) in enumerate(starts): + end = starts[idx + 1][0] if idx + 1 < len(starts) else len(lines) + body_lines = lines[line_i + 1 : end] + # Drop leading blank lines; keep rest (incl. prompts inside body — parsers skip them) + while body_lines and not body_lines[0].strip(): + body_lines = body_lines[1:] + # Trim trailing blank + while body_lines and not body_lines[-1].strip(): + body_lines = body_lines[:-1] + body = "\n".join(body_lines) + out.append( + LogSegment( + command=cmd, + body=body, + source_file=str(source_file or ""), + line_start=line_i + 1, + ) + ) + return out + + +def _decode_bytes(data: bytes) -> str: + for enc in ("utf-8", "utf-8-sig", "gb18030", "latin-1"): + try: + return data.decode(enc) + except UnicodeDecodeError: + continue + return data.decode("utf-8", errors="replace") + + +def _is_text_member(name: str) -> bool: + n = str(name or "").replace("\\", "/").strip() + if not n or n.endswith("/"): + return False + base = Path(n).name + if base.startswith(".") or base.startswith("__MACOSX"): + return False + suf = Path(base).suffix.lower() + if suf in _TEXT_SUFFIXES: + return True + # Extensionless small dumps sometimes appear; allow if no suffix + return suf == "" + + +def unpack_upload( + *, + filename: str, + data: bytes, + vendor_key: str = "", +) -> tuple[list[LogSegment], dict[str, int]]: + """Unpack a bare text upload or zip into ordered LogSegments. + + Returns (segments, stats) where stats has files / bytes / segments counts. + Raises ValueError on size / format / zip-bomb limits. + """ + name = str(filename or "upload.bin").strip() or "upload.bin" + blob = data or b"" + max_b = import_max_bytes() + max_f = import_max_files() + if len(blob) > max_b: + raise ValueError(f"upload exceeds max size {max_b}B") + + lower = name.lower() + segments: list[LogSegment] = [] + total_bytes = 0 + file_count = 0 + + if lower.endswith(".zip"): + try: + zf = zipfile.ZipFile(io.BytesIO(blob)) + except zipfile.BadZipFile as exc: + raise ValueError(f"invalid zip: {exc}") from exc + with zf: + members = [ + info + for info in zf.infolist() + if not info.is_dir() and _is_text_member(info.filename) + ] + if len(members) > max_f: + raise ValueError(f"zip has too many text files (>{max_f})") + for info in sorted(members, key=lambda x: x.filename.lower()): + if info.file_size > max_b: + raise ValueError( + f"zip member {info.filename!r} exceeds max size {max_b}B" + ) + # Zip bomb: compressed ratio / total uncompressed + total_bytes += int(info.file_size or 0) + if total_bytes > max_b: + raise ValueError(f"zip uncompressed total exceeds max size {max_b}B") + raw = zf.read(info) + text = _decode_bytes(raw) + file_count += 1 + segs = split_log_text( + text, vendor_key=vendor_key, source_file=info.filename + ) + segments.extend(segs) + else: + total_bytes = len(blob) + text = _decode_bytes(blob) + file_count = 1 + segments = split_log_text(text, vendor_key=vendor_key, source_file=name) + + return segments, { + "files": file_count, + "bytes": total_bytes, + "segments": len(segments), + } + + +def read_upload_stream(fh: BinaryIO, *, max_bytes: int | None = None) -> bytes: + """Read upload stream with a hard byte cap.""" + cap = int(max_bytes if max_bytes is not None else import_max_bytes()) + buf = fh.read(cap + 1) + if buf is None: + return b"" + if len(buf) > cap: + raise ValueError(f"upload exceeds max size {cap}B") + return buf diff --git a/netx_api/biz_state/service.py b/netx_api/biz_state/service.py index a41192c..f3ad60a 100644 --- a/netx_api/biz_state/service.py +++ b/netx_api/biz_state/service.py @@ -251,6 +251,10 @@ def update_task(db: Session, task_id: str, body: dict[str, Any]) -> dict[str, An task = db.get(BizStateTask, task_id) if not task: raise HTTPException(status_code=404, detail="task_not_found") + # Offline import tasks have no inventory NE / credentials — keep paused. + if str(task.source or "").strip().lower() == "import": + if "status" in body and str(body.get("status") or "").strip() == "running": + raise HTTPException(status_code=400, detail="import_offline_only") if "note" in body: task.note = str(body.get("note") or "")[:256] if "purpose" in body and body["purpose"] is not None: diff --git a/netx_api/biz_state_router.py b/netx_api/biz_state_router.py index 5055ffa..4a4f3ed 100644 --- a/netx_api/biz_state_router.py +++ b/netx_api/biz_state_router.py @@ -4,7 +4,7 @@ from __future__ import annotations from typing import Any -from fastapi import APIRouter, BackgroundTasks, Depends, Query +from fastapi import APIRouter, BackgroundTasks, Depends, File, Form, Query, UploadFile from fastapi.responses import StreamingResponse from pydantic import BaseModel, Field from sqlalchemy.orm import Session @@ -287,6 +287,153 @@ def api_collect_stop(task_id: str, db: Session = Depends(get_db)) -> dict[str, A return request_stop_collect(task_id) +@router.post("/import") +async def api_import_log_standalone( + background_tasks: BackgroundTasks, + file: UploadFile = File(...), + vendor_key: str = Form("zte"), + ne_name: str = Form(""), + note: str = Form(""), +) -> dict[str, Any]: + """Standalone offline import: create an import task (no NE) and parse the log.""" + from fastapi import HTTPException + + from .biz_state.import_runner import execute_import, start_standalone_import + from .biz_state.log_split import import_max_bytes, read_upload_stream + + fname = str(file.filename or "import.log") + try: + data = read_upload_stream(file.file, max_bytes=import_max_bytes()) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + result = start_standalone_import( + filename=fname, + data=data, + vendor_key=str(vendor_key or "").strip() or "zte", + ne_name=str(ne_name or "").strip(), + note=str(note or "").strip(), + ) + if not result.get("queued"): + return { + "ok": bool(result.get("ok", False)), + "started": False, + "queued": False, + "reason": result.get("reason") or "enqueue_failed", + "task_id": result.get("task_id") or "", + "created_task": bool(result.get("created_task")), + "segments_preview": result.get("segments_preview") + or {"files": result.get("files"), "segments": result.get("segments")}, + } + + segments = result.pop("segments", []) + bid = str(result.get("batch_id") or "") + tid = str(result.get("task_id") or "") + vk = str(result.get("vendor_key") or "zte") + vendor = str(result.get("vendor") or "") + device_type = str(result.get("device_type") or "") + background_tasks.add_task( + lambda: execute_import( + batch_id=bid, + task_id=tid, + segments=list(segments), + vendor_key=vk, + vendor=vendor, + device_type=device_type, + filename=fname, + ) + ) + return { + "ok": True, + "started": True, + "queued": True, + "batch_id": bid, + "task_id": tid, + "vendor_key": vk, + "filename": result.get("filename") or fname, + "ne_name": result.get("ne_name") or "", + "created_task": True, + "segments_preview": result.get("segments_preview") or {}, + "collect_running": True, + } + + +@router.post("/tasks/{task_id}/import") +async def api_import_log( + task_id: str, + background_tasks: BackgroundTasks, + file: UploadFile = File(...), + vendor_key: str = Form(""), + db: Session = Depends(get_db), +) -> dict[str, Any]: + """Upload a CLI transcript into an existing task and parse offline.""" + from fastapi import HTTPException + + from .biz_state.import_runner import execute_import, start_import_from_upload + from .biz_state.log_split import import_max_bytes, read_upload_stream + + task = db.get(BizStateTask, task_id) + if not task: + raise HTTPException(status_code=404, detail="task not found") + if bool(task.collect_running): + return { + "ok": True, + "started": False, + "reason": "already_collecting", + "task_id": task_id, + } + + fname = str(file.filename or "import.log") + try: + data = read_upload_stream(file.file, max_bytes=import_max_bytes()) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + result = start_import_from_upload( + task_id=task_id, + filename=fname, + data=data, + vendor_key=str(vendor_key or "").strip(), + ) + if not result.get("queued"): + return { + "ok": bool(result.get("ok", False)), + "started": False, + "queued": False, + "reason": result.get("reason") or "enqueue_failed", + "task_id": task_id, + "segments_preview": result.get("segments_preview") or result.get("files"), + } + + segments = result.pop("segments", []) + bid = str(result.get("batch_id") or "") + vk = str(result.get("vendor_key") or "") + vendor = str(result.get("vendor") or "") + device_type = str(result.get("device_type") or "") + background_tasks.add_task( + lambda: execute_import( + batch_id=bid, + task_id=task_id, + segments=list(segments), + vendor_key=vk, + vendor=vendor, + device_type=device_type, + filename=fname, + ) + ) + return { + "ok": True, + "started": True, + "queued": True, + "batch_id": bid, + "task_id": task_id, + "vendor_key": vk, + "filename": result.get("filename") or fname, + "segments_preview": result.get("segments_preview") or {}, + "collect_running": True, + } + + @router.get("/tasks/{task_id}/batches") def api_list_batches( task_id: str, limit: int = 50, db: Session = Depends(get_db) diff --git a/netx_api/config.py b/netx_api/config.py index ff6d852..d93654c 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -127,6 +127,9 @@ class Settings(BaseSettings): biz_state_persist_every_cmds: int = 8 # Cap raw_text loaded into Postgres from spool (0 = unlimited). biz_state_raw_max_bytes: int = 8 * 1024 * 1024 + # Manual log import: max upload / zip uncompressed bytes and text file count. + biz_state_import_max_bytes: int = 256 * 1024 * 1024 + biz_state_import_max_files: int = 200 # Dedicated biz_state worker process(es); general worker skips biz_state scheduler. biz_state_dedicated_workers: bool = True # Global ceiling for simultaneous running batches (across all workers). diff --git a/tests/test_biz_state_log_import.py b/tests/test_biz_state_log_import.py new file mode 100644 index 0000000..b5d2a1d --- /dev/null +++ b/tests/test_biz_state_log_import.py @@ -0,0 +1,367 @@ +"""Tests for biz_state manual log import (split + unpack + match path).""" + +from __future__ import annotations + +import io +import tempfile +import unittest +import zipfile +from pathlib import Path +from unittest.mock import patch + +from sqlalchemy import create_engine +from sqlalchemy.orm import sessionmaker +from sqlalchemy.pool import StaticPool + +from netx_api.biz_state import import_runner as imp +from netx_api.biz_state import spool as spool_mod +from netx_api.biz_state.command_match import match_command +from netx_api.biz_state.log_split import ( + split_log_text, + unpack_upload, +) +from netx_api.db import Base +from netx_api.models import BizStateBatch, BizStateBatchCommand, BizStateTask + + +_ZTE_SAMPLE = """ +PE1#show arp +IP Address MAC Address +10.0.0.1 aaaa.bbbb.cccc + +show\u00a0bgp\u00a0ipv4\u00a0unicast\u00a0neighbor\u00a0out\u00a01.1.1.1\u00a0| one-line +Routes Sent To This Neighbor: +Total number of routes: 1 +Network Next Hop +10.1.0.0/24 1.1.1.2 + +show weird-unknown-command +garbage line + +show bgp vpnv4 unicast neighbor in 2.2.2.2 | one-line +Routes Learned From This Neighbor: +Total number of routes: 0 +""" + +_HW_SAMPLE = """ +display ip routing-table +Routing Table +Destinations : 1 + +display bgp peer +BGP peer table +""" + + +class LogSplitTests(unittest.TestCase): + def test_zte_split_nbsp_and_prompt(self) -> None: + segs = split_log_text(_ZTE_SAMPLE, vendor_key="zte", source_file="a.ini") + cmds = [s.command for s in segs] + self.assertIn("show arp", cmds) + self.assertTrue(any(c.startswith("show bgp ipv4") for c in cmds)) + self.assertTrue(any("\u00a0" not in c for c in cmds)) + bgp = next(s for s in segs if "bgp ipv4" in s.command) + self.assertIn("10.1.0.0/24", bgp.body) + self.assertEqual(bgp.source_file, "a.ini") + + def test_preamble_before_first_show_discarded(self) -> None: + """Leading noise / prompts without show are dropped; no-show → empty.""" + raw = """Some banner +login ok +MDN-BCP-CN1-ZM8SP# +08:38:53 clock +garbage + +MDN-BCP-CN1-ZM8SP#show arp | one-line +The count is 1 +1.1.1.1 aaaa + +PE1#show interface brief +gei-0/0/0/1 up +""" + segs = split_log_text(raw, vendor_key="zte") + self.assertEqual(len(segs), 2) + self.assertEqual(segs[0].command, "show arp | one-line") + self.assertNotIn("banner", segs[0].body.lower()) + self.assertNotIn("garbage", segs[0].body) + self.assertIn("1.1.1.1", segs[0].body) + self.assertEqual(segs[1].command, "show interface brief") + self.assertEqual(split_log_text("hello\nworld\n", vendor_key="zte"), []) + + def test_hostname_hash_show_multi_segment(self) -> None: + """Real ZTE paste: MDN-...#show ... then another #show — auto split.""" + raw = """MDN-BCP-CN1-ZM8SP#show arp | one-line +08:38:53 Indonesia Sat Sep 19 2026 +The count is 2 +100.67.5.215 H d4c1.c893.1a90 gei-0/0/0/1.1409 + +MDN-BCP-CN1-ZM8SP#show interface brief +gei-0/0/0/1 up up + +PE2#show bgp vpnv4 unicast neighbor in 10.0.0.1 | one-line +Total number of routes: 0 +""" + segs = split_log_text(raw, vendor_key="zte", source_file="show-arp") + self.assertEqual(len(segs), 3) + self.assertEqual(segs[0].command, "show arp | one-line") + self.assertEqual(segs[1].command, "show interface brief") + self.assertTrue(segs[2].command.startswith("show bgp vpnv4")) + self.assertIn("100.67.5.215", segs[0].body) + self.assertIn("gei-0/0/0/1", segs[1].body) + hit0 = match_command(vendor_key="zte", command=segs[0].command) + self.assertIsNotNone(hit0) + assert hit0 is not None + self.assertEqual(hit0.profile.metric_id, "arp") + + def test_huawei_prefers_display(self) -> None: + segs = split_log_text(_HW_SAMPLE, vendor_key="huawei") + self.assertGreaterEqual(len(segs), 2) + self.assertTrue(all(s.command.lower().startswith("display") for s in segs)) + + def test_zip_multiple_files(self) -> None: + buf = io.BytesIO() + with zipfile.ZipFile(buf, "w") as zf: + zf.writestr("one.txt", "show arp\n1.1.1.1 aaaa\n") + zf.writestr("two.log", "show interface brief\ngei-0/0/0/1 up\n") + zf.writestr("skip.bin", b"\x00\x01\x02") + segs, stats = unpack_upload( + filename="dump.zip", data=buf.getvalue(), vendor_key="zte" + ) + self.assertEqual(stats["files"], 2) + self.assertEqual(stats["segments"], 2) + files = {s.source_file for s in segs} + self.assertEqual(files, {"one.txt", "two.log"}) + + def test_upload_size_limit(self) -> None: + with patch( + "netx_api.biz_state.log_split.settings" + ) as st: + st.biz_state_import_max_bytes = 100 + st.biz_state_import_max_files = 200 + with self.assertRaises(ValueError): + unpack_upload( + filename="big.txt", + data=b"x" * 200, + vendor_key="zte", + ) + + def test_match_known_zte_bgp(self) -> None: + segs = split_log_text(_ZTE_SAMPLE, vendor_key="zte") + hits = [] + misses = [] + for s in segs: + hit = match_command(vendor_key="zte", command=s.command) + if hit: + hits.append(hit.profile.metric_id) + else: + misses.append(s.command) + self.assertIn("bgp_route", hits) + self.assertTrue(any("weird-unknown" in c for c in misses)) + + +class ImportRunnerEnqueueTests(unittest.TestCase): + def setUp(self) -> None: + engine = create_engine( + "sqlite+pysqlite:///:memory:", + future=True, + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + TestingSession = sessionmaker( + bind=engine, autoflush=False, autocommit=False, expire_on_commit=False + ) + Base.metadata.create_all(bind=engine) + self.Session = TestingSession + self.db = TestingSession() + self.task = BizStateTask( + id="t-imp", + source="managed", + ne_id="ne1", + ne_name="PE1", + vendor="ZTE", + device_type="zte_zxros", + status="paused", + collect_running=False, + interval_sec=300, + ) + self.db.add(self.task) + self.db.commit() + self._sess_patch = patch.object(imp, "SessionLocal", TestingSession) + self._sess_patch.start() + + def tearDown(self) -> None: + self._sess_patch.stop() + self.db.close() + + def test_enqueue_mutex(self) -> None: + r1 = imp.enqueue_import("t-imp", filename="a.ini", vendor_key="zte") + self.assertTrue(r1.get("queued")) + self.assertTrue(r1.get("batch_id")) + r2 = imp.enqueue_import("t-imp", filename="b.ini", vendor_key="zte") + self.assertEqual(r2.get("reason"), "already_collecting") + self.db.expire_all() + task = self.db.get(BizStateTask, "t-imp") + assert task is not None + self.assertTrue(task.collect_running) + batch = self.db.get(BizStateBatch, r1["batch_id"]) + assert batch is not None + self.assertEqual(batch.status, "running") + self.assertIn("a.ini", batch.alias or "") + + +class ImportRunnerParseTests(unittest.TestCase): + def setUp(self) -> None: + self._tmpdir = tempfile.TemporaryDirectory() + self.root = Path(self._tmpdir.name) + engine = create_engine( + "sqlite+pysqlite:///:memory:", + future=True, + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + TestingSession = sessionmaker( + bind=engine, autoflush=False, autocommit=False, expire_on_commit=False + ) + Base.metadata.create_all(bind=engine) + self.Session = TestingSession + self.db = TestingSession() + self.task = BizStateTask( + id="t-imp2", + source="managed", + ne_id="ne1", + ne_name="PE1", + vendor="ZTE", + device_type="zte_zxros", + status="paused", + collect_running=True, + interval_sec=300, + ) + self.batch = BizStateBatch( + id="b-imp2", + task_id="t-imp2", + source="managed", + ne_id="ne1", + ne_name="PE1", + vendor="ZTE", + status="running", + message="importing:sample.ini", + alias="sample.ini", + ) + self.db.add(self.task) + self.db.add(self.batch) + self.db.commit() + + self._patches = [ + patch.object(imp, "SessionLocal", TestingSession), + patch.object(spool_mod.settings, "biz_state_spool_dir", str(self.root)), + ] + # collect_runner uses SessionLocal too for flush/finalize/finish + from netx_api.biz_state import collect_runner as runner + from netx_api.biz_state import persist_pool as pp + + self._patches.append(patch.object(runner, "SessionLocal", TestingSession)) + + class _SyncPersist: + def submit(self, batch_id, items): + runner._flush_spooled_commands(batch_id, list(items)) + + def wait_idle(self, *, timeout=None): + return True + + self._sync = _SyncPersist() + self._patches.append(patch.object(imp, "get_persist_pool", return_value=self._sync)) + self._patches.append(patch.object(pp, "get_persist_pool", return_value=self._sync)) + for p in self._patches: + p.start() + + def tearDown(self) -> None: + for p in self._patches: + p.stop() + self.db.close() + self._tmpdir.cleanup() + + def test_run_import_marks_unmatched_and_ok(self) -> None: + segs = split_log_text(_ZTE_SAMPLE, vendor_key="zte", source_file="sample.ini") + out = imp.run_import_batch( + batch_id="b-imp2", + task_id="t-imp2", + segments=segs, + vendor_key="zte", + vendor="ZTE", + device_type="zte_zxros", + filename="sample.ini", + ) + self.assertIn(out.get("status"), ("success", "partial", "failed")) + self.assertGreaterEqual(int(out.get("matched") or 0), 1) + self.assertGreaterEqual(int(out.get("unmatched") or 0), 1) + self.db.expire_all() + cmds = ( + self.db.query(BizStateBatchCommand) + .filter(BizStateBatchCommand.batch_id == "b-imp2") + .all() + ) + statuses = {str(c.parse_status or "") for c in cmds} + self.assertIn("unmatched", statuses) + batch = self.db.get(BizStateBatch, "b-imp2") + assert batch is not None + self.assertIn(batch.status, ("partial", "success", "failed")) + self.assertTrue(str(batch.message or "").strip()) + self.assertIn("imported:", (batch.message or "").lower()) + # Missing aux must not claim enrich in ok command messages + for c in cmds: + if c.parse_status == "ok": + self.assertNotIn("enrich=", (c.message or "")) + task = self.db.get(BizStateTask, "t-imp2") + assert task is not None + self.assertFalse(task.collect_running) + + +class StandaloneImportTaskTests(unittest.TestCase): + def setUp(self) -> None: + self.engine = create_engine( + "sqlite://", + connect_args={"check_same_thread": False}, + poolclass=StaticPool, + ) + Base.metadata.create_all(self.engine) + TestingSession = sessionmaker(bind=self.engine) + self.Session = TestingSession + self.db = TestingSession() + self._patch = patch.object(imp, "SessionLocal", TestingSession) + self._patch.start() + + def tearDown(self) -> None: + self._patch.stop() + self.db.close() + self.engine.dispose() + + def test_create_standalone_import_task(self) -> None: + out = imp.create_standalone_import_task( + vendor_key="zte", + ne_name="PE1-lab", + filename="show-arp.ini", + ) + self.assertTrue(out.get("ok")) + tid = str(out["task_id"]) + task = self.db.get(BizStateTask, tid) + assert task is not None + self.assertEqual(task.source, "import") + self.assertEqual(task.status, "paused") + self.assertEqual(task.ne_name, "PE1-lab") + self.assertEqual(task.vendor, "ZTE") + self.assertTrue(str(task.ne_id or "").startswith("import-")) + self.assertFalse(bool(task.ne_ip)) + + def test_create_defaults_name_from_filename(self) -> None: + out = imp.create_standalone_import_task( + vendor_key="huawei", + filename="C:/logs/display-ip.zip", + ) + task = self.db.get(BizStateTask, out["task_id"]) + assert task is not None + self.assertEqual(task.ne_name, "import:display-ip.zip") + self.assertEqual(task.vendor, "Huawei") + + +if __name__ == "__main__": + unittest.main() diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 8a2aa45..39e2828 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -176,6 +176,11 @@ const en = { title: "Business monitor", create: "Create task", createHint: "Pick any valid NE (managed / UME); search supported", + createModeDevice: "Pick device", + createModeImport: "Import log", + createImportHint: + "No NE needed: upload a session log or zip to create an Import task and parse the first batch. Re-import later on the task to append batches (same as multiple collects).", + createImportSubmit: "Create & import", created: "Task created", createFailed: "Create failed", empty: "No business-state tasks yet.", @@ -195,6 +200,19 @@ const en = { scheduleOff: "Manual", scheduleHint: "When checked, collect on interval; otherwise only Collect now.", collectNow: "Collect now", + importLog: "Import log", + importLogHint: + "Append a new batch to this task (same as one collect): full session or zip; discard noise before show/display; split on hostname#show; fail if no show (no SSH).", + importVendor: "Vendor", + importVendorGeneric: "Generic (show+display)", + importFile: "Log file", + importNeName: "Task name (optional)", + importNeNamePh: "e.g. PE1-offline", + importSubmit: "Import as new batch", + importing: "Importing", + importStarted: "Import started ({{count}} segments)", + importOfflineOnly: "Import tasks are offline-only (no SSH collect)", + sourceImport: "Import", stopCollect: "Stop collect", stopCollectOk: "Stop collect requested", stopCollectIdle: "No collect in progress", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index e0dadaf..f2259b6 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -176,6 +176,11 @@ const zh = { title: "业务监控", create: "创建任务", createHint: "从托管 / UME 等有效网元中选择(支持查找)", + createModeDevice: "选择设备", + createModeImport: "手动导入", + createImportHint: + "无网元也可建任务:上传会话日志或 zip,自动切分 show/display 并解析为首个批次。之后可在任务上多次「导入日志」追加批次(等同于多次采集)。", + createImportSubmit: "创建并导入", created: "已创建任务", createFailed: "创建失败", empty: "暂无业务监控任务。", @@ -195,6 +200,19 @@ const zh = { scheduleOff: "手动", scheduleHint: "勾选后按周期自动采集;不勾选则仅支持「立即采集」。", collectNow: "立即采集", + importLog: "导入日志", + importLogHint: + "向当前任务追加一批次(等同于一次采集):支持整段会话或 zip,丢掉 show/display 前杂讯,按 hostname#show 切分;整份无 show 则失败(无需 SSH)。", + importVendor: "厂商", + importVendorGeneric: "通用(show+display)", + importFile: "日志文件", + importNeName: "任务名称(可选)", + importNeNamePh: "例如:PE1-offline", + importSubmit: "导入为新批次", + importing: "导入中", + importStarted: "已开始导入({{count}} 段命令)", + importOfflineOnly: "导入任务仅支持离线导入,无 SSH 采集", + sourceImport: "导入", stopCollect: "停止采集", stopCollectOk: "已请求停止采集", stopCollectIdle: "当前没有进行中的采集", diff --git a/web/src/index.css b/web/src/index.css index 3327d46..eb97361 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -9918,7 +9918,9 @@ html.login-page--paused .login-page__flare { } .nm-page-panel .nm-config-modal__tabs button.is-active, -.bs-cmp-job-body .nm-config-modal__tabs button.is-active { +.bs-cmp-job-body .nm-config-modal__tabs button.is-active, +.app-heroui-modal .nm-config-modal__tabs button.is-active, +.nm-config-modal__tabs button.is-active { outline: 2px solid rgba(37, 99, 235, 0.55); outline-offset: 1px; box-shadow: 0 0 0 1px rgba(37, 99, 235, 0.25); diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index 7aacd42..d421669 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -20,6 +20,8 @@ import { bizStateGetBatch, bizStateGetBatchCommand, bizStateGetTask, + bizStateImportLog, + bizStateImportLogStandalone, bizStateListBatches, bizStateListBatchMetricRows, bizStateListProfiles, @@ -269,6 +271,7 @@ export function BizStatePage() { // create-task modal (all valid CLI targets) const [createOpen, setCreateOpen] = useState(false); + const [createMode, setCreateMode] = useState<"device" | "import">("device"); const [neKeyword, setNeKeyword] = useState(""); const debouncedNeKw = useDebouncedValue(neKeyword, 300); const [neSource, setNeSource] = useState("all"); @@ -324,6 +327,16 @@ export function BizStatePage() { const [rawLogText, setRawLogText] = useState(""); const [rawLogMeta, setRawLogMeta] = useState(""); const [rawLogCommandId, setRawLogCommandId] = useState(""); + + // Manual log import into an existing task (append batch) + const [importOpen, setImportOpen] = useState(false); + const [importTaskId, setImportTaskId] = useState(""); + const [importVendorKey, setImportVendorKey] = useState("zte"); + const [importNeName, setImportNeName] = useState(""); + const [importFile, setImportFile] = useState(null); + const [importBusy, setImportBusy] = useState(false); + const importFileRef = useRef(null); + const createImportFileRef = useRef(null); const [rawLogLines, setRawLogLines] = useState(0); const [rawLogRows, setRawLogRows] = useState(0); const [rawLogDeclared, setRawLogDeclared] = useState(0); @@ -649,15 +662,23 @@ export function BizStatePage() { const openCreate = () => { setCreateOpen(true); + setCreateMode("device"); setSelectedNe(null); setNeKeyword(""); setNeSource("all"); setNePage(1); + setImportVendorKey("zte"); + setImportNeName(""); + setImportFile(null); + if (createImportFileRef.current) createImportFileRef.current.value = ""; }; const closeCreate = () => { setCreateOpen(false); setSelectedNe(null); + setImportFile(null); + setImportNeName(""); + if (createImportFileRef.current) createImportFileRef.current.value = ""; }; const loadTask = async (id: string) => { @@ -725,6 +746,42 @@ export function BizStatePage() { } }; + /** Create offline import task + first batch (same entry as pick-NE create). */ + const createImportTask = async () => { + if (!importFile) return; + setImportBusy(true); + setBusy(true); + try { + const out = await bizStateImportLogStandalone(importFile, { + vendor_key: importVendorKey, + ne_name: importNeName.trim(), + }); + if (out && out.started === false) { + showError(String(out.reason || "") || t("common.opFailed")); + return; + } + const newTaskId = String(out.task_id || ""); + const n = Number(out.segments_preview?.segments || 0); + showOk( + n > 0 + ? t("bizState.importStarted", { count: String(n) }) + : t("bizState.created"), + ); + closeCreate(); + invalidateCutoverCache("bizCompare:"); + await refreshTasks(); + if (newTaskId) { + setTaskCollecting(newTaskId, true); + await openTask(newTaskId, "batches"); + } + } catch (e) { + showError(t("bizState.createFailed") + ": " + formatErr(e)); + } finally { + setImportBusy(false); + setBusy(false); + } + }; + const setScheduleEnabled = async (id: string, enabled: boolean) => { setBusy(true); try { @@ -918,6 +975,75 @@ export function BizStatePage() { await collectNowForTask(taskId, true); }; + const guessVendorKey = (vendor: string, deviceType = "") => { + const v = `${vendor || ""} ${deviceType || ""}`.toLowerCase(); + if (v.includes("huawei") || v.includes("vrp")) return "huawei"; + if (v.includes("cisco") || v.includes("ios")) return "cisco"; + if (v.includes("zte") || v.includes("zxros") || v.includes("zxr")) return "zte"; + return "zte"; + }; + + const openImportLog = (id: string, vendor = "", deviceType = "") => { + setImportTaskId(id); + setImportVendorKey(guessVendorKey(vendor, deviceType)); + setImportNeName(""); + setImportFile(null); + if (importFileRef.current) importFileRef.current.value = ""; + setImportOpen(true); + }; + + const submitImportLog = async () => { + if (!importTaskId || !importFile) return; + setImportBusy(true); + setTaskCollecting(importTaskId, true); + try { + const out = await bizStateImportLog(importTaskId, importFile, { + vendor_key: importVendorKey, + }); + if (out && out.started === false) { + setTaskCollecting(importTaskId, false); + const reason = String(out.reason || ""); + if (reason === "already_collecting") { + setTaskCollecting(importTaskId, true); + showOk(t("bizState.collecting")); + } else if (reason === "import_offline_only") { + showError(t("bizState.importOfflineOnly")); + return; + } else { + showError(reason || t("common.opFailed")); + return; + } + } else { + const n = Number(out.segments_preview?.segments || 0); + showOk( + n > 0 + ? t("bizState.importStarted", { count: String(n) }) + : t("bizState.importing"), + ); + } + setImportOpen(false); + if (taskId === importTaskId) { + setTaskTab("batches"); + try { + await refreshTaskProgress(importTaskId); + } catch { + /* poll will retry */ + } + } else { + try { + await refreshTasks(); + } catch { + /* ignore */ + } + } + } catch (e) { + setTaskCollecting(importTaskId, false); + showError(formatErr(e)); + } finally { + setImportBusy(false); + } + }; + const stopCollectForTask = async (id: string) => { setBusy(true); try { @@ -1363,7 +1489,9 @@ export function BizStatePage() { - {row.source || "managed"} + {row.source === "import" + ? t("bizState.sourceImport") + : row.source || "managed"} @@ -1388,11 +1516,15 @@ export function BizStatePage() { void setScheduleEnabled(row.id, e.target.checked)} /> - {row.status === "running" ? t("bizState.scheduleOn") : t("bizState.scheduleOff")} + {row.source === "import" + ? t("bizState.scheduleOff") + : row.status === "running" + ? t("bizState.scheduleOn") + : t("bizState.scheduleOff")} @@ -1426,32 +1558,46 @@ export function BizStatePage() { - {row.status === "running" ? ( + {row.source !== "import" ? ( + row.status === "running" ? ( + + ) : ( + + ) + ) : null} + {row.source !== "import" ? ( - ) : ( - - )} + ) : null} {row.collect_running || collectingIds[row.id] ? ( + -
- - - - - - - - - - - - {neItems.map((row) => { - const checked = - selectedNe?.id === row.id && neSourceOf(selectedNe) === neSourceOf(row); - return ( - - - - - - - - - ); - })} - {!neItems.length ? ( - - - + {createMode === "device" ? ( + <> +

{t("bizState.createHint")}

+
+ { + setNeKeyword(e.target.value); + setNePage(1); + }} + /> + { + setNeSource(e.target.value as NeSourceFilter); + setNePage(1); + }} + aria-label={t("bizState.colSource")} + > + + + + + {selectedNe ? ( + + {t("bizState.selectedNe")}: {selectedNe.name} ({selectedNe.ip_address}) ·{" "} + {neSourceOf(selectedNe)} + ) : null} -
-
- {t("bizState.colSource")}{t("bizState.colNe")}IP{t("bizState.colVendor")}{t("bizState.colConnect")}
- setSelectedNe(row)} - /> - - {row.source} - {row.name || "—"}{row.ip_address || "—"}{row.vendor || "—"} - - {row.connect_status || "—"} - -
-
- {neLoading ? t("common.refreshing") : t("bizState.neEmpty")} -
-
-
- + +
+ + + + + + + + + + + + {neItems.map((row) => { + const checked = + selectedNe?.id === row.id && + neSourceOf(selectedNe) === neSourceOf(row); + return ( + + + + + + + + + ); + })} + {!neItems.length ? ( + + + + ) : null} + +
+ {t("bizState.colSource")}{t("bizState.colNe")}IP{t("bizState.colVendor")}{t("bizState.colConnect")}
+ setSelectedNe(row)} + /> + + + {row.source} + + {row.name || "—"}{row.ip_address || "—"}{row.vendor || "—"} + + {row.connect_status || "—"} + +
+
+ {neLoading ? t("common.refreshing") : t("bizState.neEmpty")} +
+
+
+ + + ) : ( + <> +

{t("bizState.createImportHint")}

+ setImportVendorKey(e.target.value)} + fullWidth + > + + + + + + + + + )} - - + ) : ( + + )} + @@ -1622,29 +1848,33 @@ export function BizStatePage() { void setScheduleEnabled(taskId, e.target.checked)} /> {t("bizState.scheduleEnabled")} - {detail.status === "running" ? ( - + {detail.source !== "import" ? ( + detail.status === "running" ? ( + + ) : ( + + ) ) : ( - + {t("bizState.sourceImport")} )}