diff --git a/netx_api/biz_state/collect_runner.py b/netx_api/biz_state/collect_runner.py index 79a95be..5b03569 100644 --- a/netx_api/biz_state/collect_runner.py +++ b/netx_api/biz_state/collect_runner.py @@ -563,33 +563,48 @@ def _finish_task(task_id: str, *, error: str = "") -> None: db.close() -def dispatch_collect(task_id: str, *, manual: bool = False) -> None: +def dispatch_collect(task_id: str, *, manual: bool = False) -> dict[str, Any]: """Enqueue a collect round; run inline when this process owns execution. Dedicated worker mode: only enqueue (claim loop runs the batch). Inline / non-dedicated: enqueue then atomically promote+execute (skip if another worker already claimed the batch). + + Returns the enqueue result dict (ok / queued / reason / batch_id / …). """ from .claim import enqueue_collect result = enqueue_collect(task_id, manual=manual) if not result.get("queued"): - return + return result batch_id = str(result.get("batch_id") or "") if not batch_id: + return result + execute_enqueued_batch( + batch_id=batch_id, + task_id=str(result.get("task_id") or task_id), + ) + return result + + +def execute_enqueued_batch(*, batch_id: str, task_id: str) -> None: + """Promote+run a queued batch when this process owns inline execution.""" + bid = str(batch_id or "").strip() + if not bid: return - if _should_execute_inline(): - # Atomic queued→running; if false, scheduler/worker already owns it. - if not _try_claim_batch_for_execute(batch_id): - return - execute_claimed_batch( - batch_id=batch_id, - task_id=str(result.get("task_id") or task_id), - source="", - ne_id="", - vendor="", - device_type="", - ) + if not _should_execute_inline(): + return + # Atomic queued→running; if false, scheduler/worker already owns it. + if not _try_claim_batch_for_execute(bid): + return + execute_claimed_batch( + batch_id=bid, + task_id=str(task_id or ""), + source="", + ne_id="", + vendor="", + device_type="", + ) def _should_execute_inline() -> bool: diff --git a/netx_api/biz_state_router.py b/netx_api/biz_state_router.py index 465e95d..5055ffa 100644 --- a/netx_api/biz_state_router.py +++ b/netx_api/biz_state_router.py @@ -11,7 +11,7 @@ from sqlalchemy.orm import Session from .db import get_db from .biz_state import service as svc -from .biz_state.collect_runner import dispatch_collect +from .biz_state.collect_runner import execute_enqueued_batch from .lldp_shared import resolve_vendor_key from .models import BizStateTask @@ -246,10 +246,32 @@ def api_collect_now( if bool(task.collect_running): return {"ok": True, "started": False, "reason": "already_collecting", "task_id": task_id} tid = task_id - # Enqueue (and run inline only when this process owns collectors). - # Dedicated worker mode: BackgroundTasks only creates queued batch; workers claim. - background_tasks.add_task(lambda: dispatch_collect(tid, manual=True)) - return {"ok": True, "started": True, "queued": True, "task_id": task_id} + # Enqueue synchronously so collect_running flips before the HTTP response + # (schedule may stay paused; manual collect is still allowed). + from .biz_state.claim import enqueue_collect + + result = enqueue_collect(tid, manual=True) + 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, + } + batch_id = str(result.get("batch_id") or "") + if batch_id: + background_tasks.add_task( + lambda bid=batch_id, t=tid: execute_enqueued_batch(batch_id=bid, task_id=t) + ) + return { + "ok": True, + "started": True, + "queued": True, + "batch_id": batch_id, + "task_id": tid, + "collect_running": True, + } @router.post("/tasks/{task_id}/collect/stop") diff --git a/web/src/i18n/en.ts b/web/src/i18n/en.ts index 0cc19bb..585dbea 100644 --- a/web/src/i18n/en.ts +++ b/web/src/i18n/en.ts @@ -189,7 +189,7 @@ const en = { colConnect: "Connect", start: "Start schedule", pause: "Pause", - statusPaused: "Paused", + statusPaused: "Manual only", scheduleEnabled: "Enable schedule", scheduleOn: "Scheduled", scheduleOff: "Manual", @@ -204,6 +204,7 @@ const en = { started: "Schedule enabled", paused: "Schedule paused (manual only)", collecting: "Collecting", + statusCollecting: "Collecting", detail: "Details", interval: "Interval", intervalUnitDays: "days", diff --git a/web/src/i18n/zh.ts b/web/src/i18n/zh.ts index f1ddf2a..8e0f5cf 100644 --- a/web/src/i18n/zh.ts +++ b/web/src/i18n/zh.ts @@ -189,7 +189,7 @@ const zh = { colConnect: "连通", start: "启动周期", pause: "暂停", - statusPaused: "已暂停", + statusPaused: "仅手动", scheduleEnabled: "启用周期调度", scheduleOn: "周期", scheduleOff: "手动", @@ -204,6 +204,7 @@ const zh = { started: "已启用周期采集", paused: "已暂停周期(仅手动采集)", collecting: "采集中", + statusCollecting: "采集中", detail: "详情", interval: "采集周期", intervalUnitDays: "天", diff --git a/web/src/index.css b/web/src/index.css index 5c82c03..69ee732 100644 --- a/web/src/index.css +++ b/web/src/index.css @@ -10072,7 +10072,40 @@ html.login-page--paused .login-page__flare { display: flex; align-items: center; gap: 8px; - flex: 0 0 auto; + flex: 0 1 auto; + flex-wrap: wrap; + max-width: 100%; +} + +.bs-commands-card-list { + display: flex; + flex-direction: column; + gap: 8px; + max-height: min(48vh, 520px); + overflow: auto; + padding: 2px; +} + +.bs-commands-card { + padding: 8px 10px; + border-radius: 8px; + border: 1px solid rgba(148, 163, 184, 0.22); + background: rgba(15, 23, 42, 0.45); +} + +.bs-commands-card__msg { + flex: 1 0 100%; + margin-top: 2px; + font-size: 12px; + color: #94a3b8; + overflow: hidden; + text-overflow: ellipsis; + white-space: nowrap; +} + +.bs-cmd-metric { + font-size: 11px; + font-family: ui-monospace, SFMono-Regular, Menlo, Consolas, monospace; } .bs-inline-cmd { diff --git a/web/src/pages/network/BizStatePage.tsx b/web/src/pages/network/BizStatePage.tsx index 012b917..debd54a 100644 --- a/web/src/pages/network/BizStatePage.tsx +++ b/web/src/pages/network/BizStatePage.tsx @@ -467,13 +467,13 @@ export function BizStatePage() { const commandsSheetColumns = useMemo( () => [ + { key: "_actions", header: t("bizState.colActions") }, { key: "raw_command", header: t("bizState.colCommand") }, { key: "metric_id", header: "metric" }, { key: "parse_status", header: t("bizState.colStatus") }, { key: "raw_line_count", header: t("bizState.colRawLines") }, { key: "row_count", header: t("bizState.colRows") }, { key: "message", header: t("bizState.colMessage") }, - { key: "_actions", header: t("bizState.colActions") }, ], [t], ); @@ -823,8 +823,20 @@ export function BizStatePage() { const collectNowForTask = async (id: string, fromModal = false) => { setTaskCollecting(id, true); try { - await bizStateCollectNow(id); - showOk(t("bizState.collecting")); + const out = await bizStateCollectNow(id); + if (out && out.started === false) { + setTaskCollecting(id, false); + const reason = String(out.reason || ""); + if (reason === "already_collecting") { + setTaskCollecting(id, true); + showOk(t("bizState.collecting")); + } else { + showError(reason || t("common.opFailed")); + return; + } + } else { + showOk(t("bizState.collecting")); + } if (fromModal && taskId === id) { setTaskTab("batches"); try { @@ -1295,8 +1307,8 @@ export function BizStatePage() {
- {row.collect_running ? ( - {t("bizState.collecting")} + {row.collect_running || collectingIds[row.id] ? ( + {t("bizState.statusCollecting")} ) : row.status === "running" ? ( {t("bizState.scheduleOn")} ) : row.status === "paused" ? ( @@ -1500,13 +1512,21 @@ export function BizStatePage() { {detail?.ne_name || detail?.ne_ip || t("bizState.detail")} - {detail ? ` · ${detail.status}` : ""} {detail ? (
+ {detail.collect_running || (taskId && collectingIds[taskId]) ? ( + {t("bizState.statusCollecting")} + ) : detail.status === "running" ? ( + {t("bizState.scheduleOn")} + ) : detail.status === "paused" ? ( + {t("bizState.statusPaused")} + ) : ( + {t("bizState.scheduleOff")} + )}
))} @@ -2238,6 +2270,64 @@ export function BizStatePage() {
) : null} + {activeSheet.id === "commands" ? ( +
+ {displayRows.map((row, i) => ( +
+ + {cellText(row.raw_command) || "—"} + +
+ + + + {t("bizState.rawLogStats", { + lines: Number(row.raw_line_count ?? 0), + rows: Number(row.row_count ?? 0), + })} + + {cellText(row.parse_status) ? ( + + {cellText(row.parse_status)} + + ) : null} + {cellText(row.metric_id) ? ( + {cellText(row.metric_id)} + ) : null} +
+ {cellText(row.message) ? ( +
+ {cellText(row.message)} +
+ ) : null} +
+ ))} + {!displayRows.length ? ( +
{t("bizState.sheetEmpty")}
+ ) : null} +
+ ) : null} + + {activeSheet.id !== "commands" ? (
+ ) : ( +
+ { + setSheetKeyword(e.target.value); + setSheetPage(1); + }} + /> + + {`${displayTotal} ${t("bizState.colRows")}`} + +
+ )} + {activeSheet.id !== "commands" ? (
@@ -2327,7 +2433,7 @@ export function BizStatePage() { ) : null} - {sheetLoading && activeSheet.id !== "commands" && !displayRows.length ? ( + {sheetLoading && !displayRows.length ? (
{t("bizState.sheetLoading")}
@@ -2337,6 +2443,7 @@ export function BizStatePage() {
+ ) : null} apiDelete<{ ok: boolean }>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}`); export const bizStateCollectNow = (taskId: string) => - apiPost<{ ok: boolean }>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}/collect`, {}); + apiPost<{ + ok: boolean; + started?: boolean; + queued?: boolean; + reason?: string; + batch_id?: string; + collect_running?: boolean; + }>(`/v1/biz-state/tasks/${encodeURIComponent(taskId)}/collect`, {}); export const bizStateCollectStop = (taskId: string) => apiPost<{