diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index b9219ec..b17569a 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -407,6 +407,7 @@ def _run_collect_lane( if holder.get("timed_out"): raise TimeoutError(f"{label}_aborted") cmd_count += 1 + # Persist "running" before CLI so UI shows cmd progress during long reads. cmd_row = BizStateBatchCommand( id=uuid4().hex, batch_id=batch_id, @@ -414,8 +415,14 @@ def _run_collect_lane( profile_id=profile_id, raw_command=concrete[:512], params_json=dict(params or {}), + parse_status="running", + message="collecting", created_at=_utcnow(), ) + sdb.add(cmd_row) + sdb.commit() + _bump_batch_progress(batch_id, add_cmds=1) + cache_hit_primary = False try: cached = session.get_cached(concrete) @@ -494,6 +501,7 @@ def _run_collect_lane( cmd_count += 1 sdb.add(aux_row) sdb.commit() + _bump_batch_progress(batch_id, add_cmds=1) continue resolved_aux.append(ra) aux_row = BizStateBatchCommand( @@ -505,9 +513,14 @@ def _run_collect_lane( metric_id="", raw_command=ra.command[:512], params_json={}, + parse_status="running", + message=f"aux_for={cmd_row.id};collecting"[:1020], created_at=_utcnow(), ) cmd_count += 1 + sdb.add(aux_row) + sdb.commit() + _bump_batch_progress(batch_id, add_cmds=1) entry, cache_hit = session.fetch_and_parse( ra.command, parser_id=ra.parser_id, @@ -610,6 +623,8 @@ def _run_collect_lane( any_ok = True sdb.add(cmd_row) sdb.commit() + if n: + _bump_batch_progress(batch_id, add_rows=n) finally: sdb.close() return total_rows, cmd_count, any_fail, any_ok @@ -650,6 +665,37 @@ def _absorb_lane_result( ) +def _bump_batch_progress(batch_id: str, *, add_cmds: int = 0, add_rows: int = 0) -> None: + """Atomically bump batch counters so UI can show progress while lanes still run.""" + cmds = int(add_cmds or 0) + rows = int(add_rows or 0) + if not batch_id or (cmds <= 0 and rows <= 0): + return + db = SessionLocal() + try: + batch = ( + db.query(BizStateBatch) + .filter(BizStateBatch.id == batch_id) + .with_for_update() + .one_or_none() + ) + if not batch: + return + if cmds > 0: + batch.command_count = int(batch.command_count or 0) + cmds + if rows > 0: + batch.row_count = int(batch.row_count or 0) + rows + db.commit() + except Exception: + _log.exception("biz_state bump batch progress failed batch=%s", batch_id) + try: + db.rollback() + except Exception: + pass + finally: + db.close() + + def _run_collect_session( *, task_id: str, @@ -844,12 +890,17 @@ def _run_collect_session( ) if lane_errors and not any_ok and cmd_count == 0: - raise RuntimeError("; ".join(lane_errors)[:1020]) + # Progressive bumps may already have cmds; only hard-fail if nothing landed. + live = db.get(BizStateBatch, batch_id) + if not live or (int(live.command_count or 0) == 0 and int(live.row_count or 0) == 0): + raise RuntimeError("; ".join(lane_errors)[:1020]) + any_fail = True batch = db.get(BizStateBatch, batch_id) if batch: - batch.command_count = cmd_count - batch.row_count = total_rows + # Prefer progressive counters (survive lane timeout) over in-memory lane totals. + batch.command_count = max(int(batch.command_count or 0), int(cmd_count or 0)) + batch.row_count = max(int(batch.row_count or 0), int(total_rows or 0)) batch.ended_at = _utcnow() if any_fail and any_ok: batch.status = "partial" diff --git a/netx_api/config.py b/netx_api/config.py index 2b7b70c..3354c0c 100644 --- a/netx_api/config.py +++ b/netx_api/config.py @@ -108,8 +108,9 @@ class Settings(BaseSettings): biz_state_scheduler_tick_sec: int = 15 biz_state_dispatch_workers: int = 2 # Heavy CLI lane (interface detail / routes / MAC): own SSH + longer timeouts. - biz_state_heavy_read_timeout_sec: int = 300 - biz_state_heavy_run_timeout_cap_sec: int = 900 + # show interface on large boxes can take ~20 minutes for a single command. + biz_state_heavy_read_timeout_sec: int = 1500 + biz_state_heavy_run_timeout_cap_sec: int = 2400 biz_state_heavy_workers: int = 4 # Managed NE exec: max CLI commands per request (lab can raise; hard-capped in ne_exec). ne_exec_max_commands: int = 5 diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index 529355e..8e9739f 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -270,6 +270,27 @@ export function BizStatePage() { // eslint-disable-next-line react-hooks/exhaustive-deps -- refresh when purpose filter / refreshTasks changes }, [refreshTasks]); + // While a collect is running, refresh batch counters so cmd/row progress is visible. + useEffect(() => { + if (!taskId || !detail?.collect_running) return; + let cancelled = false; + const tick = async () => { + try { + if (cancelled) return; + await loadTask(taskId); + } catch { + /* ignore transient poll errors */ + } + }; + const timer = window.setInterval(() => void tick(), 3000); + void tick(); + return () => { + cancelled = true; + window.clearInterval(timer); + }; + // eslint-disable-next-line react-hooks/exhaustive-deps -- poll while collect_running + }, [taskId, detail?.collect_running]); + useEffect(() => { if (!createOpen) return; let cancelled = false; @@ -640,13 +661,18 @@ export function BizStatePage() { try { await bizStateCollectNow(id); showOk(t("bizState.collecting")); - for (let i = 0; i < 20; i++) { - await new Promise((r) => setTimeout(r, 1500)); + if (fromModal && taskId === id) { + setTaskTab("batches"); + } + // Heavy show-interface can take ~20 minutes; poll long enough and refresh batches. + const deadline = Date.now() + 32 * 60 * 1000; + while (Date.now() < deadline) { + await new Promise((r) => setTimeout(r, 3000)); if (fromModal && taskId === id) { try { + await loadTask(id); const task = await bizStateGetTask(id); setDetail(task); - await refreshTasks(); if (!task.collect_running) break; } catch { break; @@ -660,7 +686,6 @@ export function BizStatePage() { await refreshTasks(); if (fromModal && taskId === id) { await loadTask(id); - setTaskTab("batches"); } } catch (e) { showError(formatErr(e));