This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/dev/pr-12398-a93005261d6f0af1a6a4f9d250fed632d34c6f86 in repository https://gitbox.apache.org/repos/asf/seatunnel.git
commit bc39c16cb0980b9f6ee310f54dc62efe809f1e7d Author: Goutam Adwant <[email protected]> AuthorDate: Mon Sep 28 15:05:09 2026 +0000 [Feature][Zeta] Show worker resource snapshots in the Web UI (#12398) Signed-off-by: Goutam Adwant <[email protected]> --- docs/en/engines/zeta/web-ui.md | 48 ++- docs/images/ui/workers-details.png | Bin 0 -> 117872 bytes docs/images/ui/workers-narrow.png | Bin 0 -> 69300 bytes docs/images/ui/workers.png | Bin 75261 -> 110920 bytes docs/zh/engines/zeta/web-ui.md | 38 +- seatunnel-engine/seatunnel-engine-ui/README.md | 5 + .../seatunnel-engine-ui/cypress/e2e/managers.cy.ts | 243 +++++++++++++ .../src/locales/en_US/managers.ts | 30 +- .../src/locales/zh_CN/managers.ts | 29 +- .../src/service/manager/index.ts | 6 +- .../src/service/manager/types.ts | 23 ++ .../seatunnel-engine-ui/src/tests/managers.spec.ts | 388 +++++++++++++++++++-- .../src/views/managers/index.tsx | 306 +++++++++++++--- .../src/views/managers/resources.ts | 85 +++++ 14 files changed, 1103 insertions(+), 98 deletions(-) diff --git a/docs/en/engines/zeta/web-ui.md b/docs/en/engines/zeta/web-ui.md index 6b49006e57..8f39fabad9 100644 --- a/docs/en/engines/zeta/web-ui.md +++ b/docs/en/engines/zeta/web-ui.md @@ -92,7 +92,53 @@ The "Finished Jobs" section displays jobs that have reached a terminal state, su The "Workers" section displays system monitoring information for worker nodes. Use it to inspect worker address, resource status, and runtime health signals exposed by the engine. - +The table shows process CPU, heap used/max, physical memory, GC counts, threads, +and slots. **Details** opens all system monitoring fields and the worker's +resource-manager snapshot: available/total CPU and heap resources, heartbeat +CPU/memory usage, tags, and running job count. + +The worker table scrolls horizontally on narrow screens. The Details column +scrolls with the data instead of covering it, and long slot descriptions wrap +within their column. The existing sidebar collapse control remains available. + +- Fixed-slot workers show used/total and free slots. Dynamic-slot workers show + only used slots and an explicit dynamic label: tracked slots are not capacity. +- Missing values are shown as `—`, not zero. Monitoring-only and resource-only + workers remain visible; an unavailable endpoint displays a warning and clears + its old values. An unavailable resource snapshot is not an empty cluster. +- The page refreshes 30 seconds after the previous requests finish, with at most + one refresh in flight. **Refresh** requests an immediate update. Polling pauses + while the browser tab is hidden and refreshes when it becomes visible again. + An already-running refresh is allowed to finish before a new one starts. + Leaving the page stops polling. The table paginates locally; each Workers + refresh sends two HTTP requests from the browser. On the server, the monitoring + endpoint dispatches one RPC per cluster member concurrently, so its fan-out + is O(n) for n members, not constant-cost. Responses are collected against one + shared deadline (`seatunnel.engine.health-metrics-timeout-seconds`, 3 seconds + by default); members that miss it are reported with a `timeout` error marker. + The browser's 6-second timeout does not cancel server-side operations. This + UI change does not alter backend RPC or timeout behavior. +- Monitoring and resource-manager values are separate samples. **Resource + response time** is when the master built the resource response, not when a + worker last sent a heartbeat. It cannot establish heartbeat freshness. + +This is a read-only view using the existing monitoring and +[`/resource/workers`](./rest-api-v2.md) endpoints. Task-to-worker drill-down and +historical metrics are not included. The Master page shows monitoring details +only and does not request worker resource data. + +The screenshots below show the actual UI with deterministic Cypress REST fixtures, +not a live cluster. The table includes a fixed-slot worker and a dynamic-slot worker; +missing measurements remain unavailable rather than appearing as zero. + + + + + +On a narrow screen, collapse the sidebar and scroll the table horizontally to +inspect the slot summary or reach Details. + + ## Master diff --git a/docs/images/ui/workers-details.png b/docs/images/ui/workers-details.png new file mode 100644 index 0000000000..894e054b76 Binary files /dev/null and b/docs/images/ui/workers-details.png differ diff --git a/docs/images/ui/workers-narrow.png b/docs/images/ui/workers-narrow.png new file mode 100644 index 0000000000..1f86c93f3f Binary files /dev/null and b/docs/images/ui/workers-narrow.png differ diff --git a/docs/images/ui/workers.png b/docs/images/ui/workers.png index a2bf39ec21..dca677aadd 100644 Binary files a/docs/images/ui/workers.png and b/docs/images/ui/workers.png differ diff --git a/docs/zh/engines/zeta/web-ui.md b/docs/zh/engines/zeta/web-ui.md index 0db8407644..4d1cb43699 100644 --- a/docs/zh/engines/zeta/web-ui.md +++ b/docs/zh/engines/zeta/web-ui.md @@ -93,7 +93,43 @@ Web UI 不负责提交作业,也不提供 cancel、stop、savepoint、restore “工作节点”模块展示 worker 节点的系统监控信息。可以用它查看 worker 地址、资源状态和引擎暴露的运行时健康信号。 - +表格展示进程 CPU、堆内存已用/上限、物理内存、GC 次数、线程数和槽位。 +点击**详情**可查看全部系统监控字段,以及资源管理器快照中的可用/总计 CPU +和堆内存资源、心跳 CPU/内存使用率、标签和运行中作业数。 + +窄屏下可横向滚动工作节点表格。详情列随数据一起滚动,不再遮挡数据,较长的 +槽位说明会在列内换行。仍可使用现有的侧边栏折叠控件。 + +- 固定槽位 Worker 展示已用/总计及空闲槽位;动态槽位 Worker 仅展示已用 + 槽位并标注“动态”,已跟踪槽位数量并不代表容量。 +- 缺失值显示为 `—`,而非零。只有监控或只有资源快照的 Worker 仍会显示; + 接口不可用时会显示警告并清除旧值。资源快照不可用不代表集群为空。 +- 上次请求完成后每 30 秒刷新一次,同时最多执行一轮刷新。点击**刷新** + 可立即更新。浏览器标签页隐藏时暂停轮询,恢复可见时重新刷新;若已有 + 请求正在执行,则等待该轮完成后再开始新一轮。离开页面后停止轮询。 + 表格在客户端分页,Worker 页面每轮刷新由浏览器发送两个 HTTP 请求。 + 服务端监控接口会并发向每个集群成员发送一次 RPC,因此 n 个成员 + 对应 O(n) 次 RPC,而非固定成本。响应使用同一个截止时间收集 + (`seatunnel.engine.health-metrics-timeout-seconds`,默认 3 秒); + 超时的成员会以带有 `timeout` 错误标记的条目返回。浏览器的 6 秒 + 超时不会取消服务端操作。此 UI 变更不修改后端 RPC 或超时行为。 +- 监控与资源管理器数据为独立采样。**资源响应时间**为 Master 构建响应的 + 时间,并非 Worker 最近一次心跳的时间,不能用于判断心跳新鲜度。 + +这是基于现有系统监控和 [`/resource/workers`](./rest-api-v2.md) 接口的只读视图。 +暂不包含任务与 Worker 的关联下钻或历史指标。管理节点页面仅展示系统监控 +详情,不请求 Worker 资源数据。 + +以下截图展示实际 UI,使用确定性的 Cypress REST 测试数据,并非真实集群。 +表格包含固定槽位和动态槽位 Worker;缺失的监控值显示为不可用,而不是零。 + + + + + +在窄屏上,可以收起侧边栏并水平滚动表格,查看槽位信息或访问详情按钮。 + + ## 管理节点 diff --git a/seatunnel-engine/seatunnel-engine-ui/README.md b/seatunnel-engine/seatunnel-engine-ui/README.md index 1d80a8f4bf..324233b278 100644 --- a/seatunnel-engine/seatunnel-engine-ui/README.md +++ b/seatunnel-engine/seatunnel-engine-ui/README.md @@ -39,6 +39,11 @@ npm run test:unit ### Run End-to-End Tests with [Cypress] +The Cypress specs, including the worker resource fixtures, are manual/optional +checks. The `seatunnel-ui` job in `.github/workflows/backend.yml` runs lint, unit +tests and the build, but does not run Cypress. These REST-contract fixtures do +not validate a deployed cluster. + ```sh npm run test:e2e:dev ``` diff --git a/seatunnel-engine/seatunnel-engine-ui/cypress/e2e/managers.cy.ts b/seatunnel-engine/seatunnel-engine-ui/cypress/e2e/managers.cy.ts new file mode 100644 index 0000000000..29f5baf5fb --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-ui/cypress/e2e/managers.cy.ts @@ -0,0 +1,243 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +describe('Worker resources', () => { + beforeEach(() => { + cy.intercept('GET', '**/overview', { + projectVersion: 'test', + gitCommitAbbrev: 'test', + totalSlot: '4', + unassignedSlot: '1', + workers: '2', + runningJobs: '1' + }) + cy.intercept('GET', '**/system-monitoring-information', [ + { + host: 'fixed', + port: '5801', + isMaster: 'false', + 'load.process': '10%', + 'heap.memory.used': '1G', + 'heap.memory.max': '4G', + 'thread.count': '12' + }, + { host: '2001:db8::1', port: '5801', isMaster: 'false' }, + { host: 'master', port: '5801', isMaster: 'true' } + ]).as('monitoring') + cy.intercept('GET', '**/resource/workers', { + available: true, + collectedAt: 1723017600000, + workers: [ + { + address: '[fixed]:5801', + dynamicSlot: false, + totalSlots: 4, + usedSlots: 3, + freeSlots: 1, + availableCpuCores: 0, + totalCpuCores: 8, + cpuUsage: 0, + tags: { region: '<img src=x onerror=alert(1)>' }, + runningJobIds: [] + }, + { + address: '[2001:db8::1]:5801', + dynamicSlot: true, + totalSlots: 999, + usedSlots: 2, + freeSlots: 997 + } + ] + }).as('resources') + }) + + it('shows fixed and dynamic slots and read-only details', () => { + cy.visit('/#/managers/workers') + cy.wait(['@monitoring', '@resources']) + cy.contains('3 / 4 used, 1 free') + cy.contains('Dynamic — 2 used (no fixed capacity)') + cy.contains('999').should('not.exist') + cy.contains('td', 'fixed:5801').parent().contains('button', 'Details').click() + cy.get('.n-drawer') + .should('be.visible') + .within(() => { + cy.contains('0 / 8') + cy.contains('0.0%') + cy.contains('heap.memory.max') + cy.contains('<img src=x onerror=alert(1)>') + cy.get('img').should('not.exist') + }) + }) + + it('clears previous slots on an unavailable snapshot and recovers', () => { + cy.visit('/#/managers/workers') + cy.wait(['@monitoring', '@resources']) + cy.contains('3 / 4 used, 1 free') + cy.intercept('GET', '**/resource/workers', { + available: false, + collectedAt: 1723017610000, + workers: [] + }).as('unavailable') + cy.contains('button', 'Refresh').click() + cy.wait('@unavailable') + cy.contains('Worker resources are unavailable') + cy.contains('3 / 4 used, 1 free').should('not.exist') + cy.contains('fixed:5801') + cy.intercept('GET', '**/resource/workers', { + available: true, + collectedAt: 1723017620000, + workers: [ + { address: 'fixed:5801', dynamicSlot: false, totalSlots: 4, usedSlots: 0, freeSlots: 4 } + ] + }).as('recovered') + cy.contains('button', 'Refresh').click() + cy.wait('@recovered') + cy.contains('0 / 4 used, 4 free') + cy.contains('Worker resources are unavailable').should('not.exist') + }) + + it('pauses hidden-tab polling and refreshes when visible again', () => { + cy.clock(0, ['setTimeout', 'clearTimeout']) + cy.visit('/#/managers/workers', { + onBeforeLoad(win) { + const requests = cy.spy(win.XMLHttpRequest.prototype, 'open') + requests + .withArgs('GET', Cypress.sinon.match(/\/system-monitoring-information$/)) + .as('monitoringStarts') + requests.withArgs('GET', Cypress.sinon.match(/\/resource\/workers$/)).as('resourceStarts') + } + }) + cy.wait(['@monitoring', '@resources']) + cy.contains('button', 'Refresh').should('not.be.disabled') + // Prove the real page's polling timer is armed before exercising the pause. + cy.tick(30_000) + cy.wait(['@monitoring', '@resources']) + cy.contains('button', 'Refresh').should('not.be.disabled') + cy.get('@monitoringStarts').should('have.been.calledTwice') + cy.get('@resourceStarts').should('have.been.calledTwice') + cy.document().then((doc) => { + Object.defineProperty(doc, 'visibilityState', { configurable: true, get: () => 'hidden' }) + doc.dispatchEvent(new Event('visibilitychange')) + }) + cy.tick(90_000) + // Let timer-triggered Axios promise chains drain before reading page-side calls. + // Unlike proxy alias counts, these observe request starts without network latency. + cy.window().then( + (win) => + new Cypress.Promise<void>((resolve) => { + const channel = new win.MessageChannel() + channel.port1.onmessage = () => { + channel.port1.close() + channel.port2.close() + resolve() + } + channel.port2.postMessage(null) + }) + ) + cy.get('@monitoringStarts').should('have.been.calledTwice') + cy.get('@resourceStarts').should('have.been.calledTwice') + cy.document().then((doc) => { + Object.defineProperty(doc, 'visibilityState', { configurable: true, get: () => 'visible' }) + doc.dispatchEvent(new Event('visibilitychange')) + }) + cy.wait(['@monitoring', '@resources']) + cy.contains('button', 'Refresh').should('not.be.disabled') + cy.get('@monitoringStarts').should('have.been.calledThrice') + cy.get('@resourceStarts').should('have.been.calledThrice') + cy.contains('3 / 4 used, 1 free') + }) + + it('does not fetch worker slots on the Master page', () => { + cy.visit('/#/managers/master') + cy.wait('@monitoring') + cy.contains('master:5801') + cy.contains('th', 'Slots').should('not.exist') + cy.get('@resources.all').should('have.length', 0) + }) + + for (const width of [1440, 768, 390]) { + it(`keeps worker cells readable while scrolling at ${width}px`, () => { + cy.viewport(width, 900) + cy.visit('/#/managers/workers') + cy.wait(['@monitoring', '@resources']) + cy.get('.n-data-table .n-scrollbar-container').as('tableViewport') + cy.get('@tableViewport').should(($viewport) => { + expect($viewport[0].scrollWidth).to.be.greaterThan($viewport[0].clientWidth) + }) + // At the left edge an off-screen action must not cover the address column. + cy.get('@tableViewport').scrollTo('left') + cy.contains('td', 'fixed:5801').should('be.visible') + cy.get('.n-data-table tbody tr') + .first() + .find('td') + .last() + .should(($action) => { + const action = $action[0] + const viewport = action.closest('.n-scrollbar-container')! + expect(action.getBoundingClientRect().right).to.be.greaterThan( + viewport.getBoundingClientRect().right + ) + }) + cy.get('@tableViewport').scrollTo('right') + cy.contains('td', 'Dynamic — 2 used (no fixed capacity)').should(($cell) => { + const cell = $cell[0] + const action = cell.nextElementSibling! + const bounds = cell.getBoundingClientRect() + const range = cell.ownerDocument.createRange() + range.selectNodeContents(cell.firstElementChild || cell) + const lines = Array.from(range.getClientRects()).filter((rect) => rect.width > 0) + expect(lines.length, 'wrapped slot text').to.be.greaterThan(1) + for (const line of lines) { + expect(line.left, 'text stays inside its cell').to.be.at.least(bounds.left) + expect(line.right, 'text stays before the next cell').to.be.at.most(bounds.right) + } + expect( + action.getBoundingClientRect().left, + 'action does not overlap slot cell' + ).to.be.at.least(bounds.right - 1) + }) + cy.get('.n-data-table tbody tr').first().contains('button', 'Details').click() + cy.get('.n-drawer').should('be.visible').and('contain', 'fixed:5801') + if (width === 390) { + cy.get('.n-drawer-header__close').click() + cy.get('.n-drawer').should('not.exist') + cy.get('.n-layout-toggle-bar').click() + cy.get('.n-layout-sider').should(($sidebar) => { + expect($sidebar[0].getBoundingClientRect().width).to.be.at.most(65) + }) + cy.contains('td', 'Dynamic — 2 used (no fixed capacity)').scrollIntoView() + cy.contains('td', 'Dynamic — 2 used (no fixed capacity)').then(($cell) => { + cy.get('@tableViewport').scrollTo($cell[0].offsetLeft, 0) + }) + cy.contains('td', 'Dynamic — 2 used (no fixed capacity)').should(($cell) => { + const cell = $cell[0] + const viewport = cell.closest('.n-scrollbar-container')!.getBoundingClientRect() + const range = cell.ownerDocument.createRange() + range.selectNodeContents(cell.firstElementChild || cell) + for (const line of Array.from(range.getClientRects())) { + expect(line.left, 'slot text visible after sidebar collapse').to.be.at.least( + viewport.left + ) + expect(line.right, 'slot text visible after sidebar collapse').to.be.at.most( + viewport.right + ) + } + }) + } + }) + } +}) diff --git a/seatunnel-engine/seatunnel-engine-ui/src/locales/en_US/managers.ts b/seatunnel-engine/seatunnel-engine-ui/src/locales/en_US/managers.ts index e2e4d7ccdf..c0cbea685e 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/locales/en_US/managers.ts +++ b/seatunnel-engine/seatunnel-engine-ui/src/locales/en_US/managers.ts @@ -16,5 +16,33 @@ */ export default { - managers: 'Managers' + managers: 'Managers', + address: 'Address', + cpu: 'Process CPU', + heap: 'Heap used / max', + physical: 'Physical memory', + gc: 'GC count (minor / major)', + threads: 'Threads', + slots: 'Slots', + details: 'Details', + refresh: 'Refresh', + refresh_hint: 'Refreshes every 30 seconds after the previous request completes.', + monitor_unavailable: 'System monitoring is unavailable. Previous values have been cleared.', + resource_unavailable: + 'Worker resources are unavailable. This does not mean the cluster has no workers.', + snapshot_hint: + 'Monitoring and worker resources are independent snapshots. Resource values come from the latest heartbeat; response time is not heartbeat freshness.', + collected_at: 'Resource response time', + dynamic_used: 'Dynamic — {used} used (no fixed capacity)', + fixed_slots: '{used} / {total} used, {free} free', + cpu_resources: 'CPU cores (available / total)', + heap_resources: 'Heap resources (available / total)', + cpu_usage: 'Heartbeat CPU usage', + memory_usage: 'Heartbeat memory usage', + running_jobs: 'Running job count', + tags: 'Tags', + resources: 'Worker resources', + monitoring: 'System monitoring', + resource_missing: 'No resource snapshot is available for this worker.', + monitor_missing: 'No system monitoring sample is available for this worker.' } diff --git a/seatunnel-engine/seatunnel-engine-ui/src/locales/zh_CN/managers.ts b/seatunnel-engine/seatunnel-engine-ui/src/locales/zh_CN/managers.ts index cd924b0d05..dafbf9b156 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/locales/zh_CN/managers.ts +++ b/seatunnel-engine/seatunnel-engine-ui/src/locales/zh_CN/managers.ts @@ -16,5 +16,32 @@ */ export default { - managers: '管理者' + managers: '管理者', + address: '地址', + cpu: '进程 CPU', + heap: '堆内存已用 / 上限', + physical: '物理内存', + gc: 'GC 次数(Minor / Major)', + threads: '线程数', + slots: '槽位', + details: '详情', + refresh: '刷新', + refresh_hint: '上次请求完成后每 30 秒刷新一次。', + monitor_unavailable: '系统监控暂不可用,已清除之前的数值。', + resource_unavailable: 'Worker 资源暂不可用,并不代表集群没有 Worker。', + snapshot_hint: + '系统监控与 Worker 资源为独立快照。资源值来自最近一次心跳;响应时间不表示心跳新鲜度。', + collected_at: '资源响应时间', + dynamic_used: '动态 — 已用 {used}(无固定容量)', + fixed_slots: '已用 {used} / {total},空闲 {free}', + cpu_resources: 'CPU 核心(可用 / 总计)', + heap_resources: '堆内存资源(可用 / 总计)', + cpu_usage: '心跳 CPU 使用率', + memory_usage: '心跳内存使用率', + running_jobs: '运行中作业数', + tags: '标签', + resources: 'Worker 资源', + monitoring: '系统监控', + resource_missing: '该 Worker 暂无资源快照。', + monitor_missing: '该 Worker 暂无系统监控采样。' } diff --git a/seatunnel-engine/seatunnel-engine-ui/src/service/manager/index.ts b/seatunnel-engine/seatunnel-engine-ui/src/service/manager/index.ts index 777b853cf4..f44696d986 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/service/manager/index.ts +++ b/seatunnel-engine/seatunnel-engine-ui/src/service/manager/index.ts @@ -16,9 +16,11 @@ */ import { get } from '@/service/service' -import type { Monitor } from './types' +import type { Monitor, WorkerResourceSnapshot } from './types' export const getMonitors = () => get<Monitor[]>('/system-monitoring-information') +export const getWorkerResources = () => get<WorkerResourceSnapshot>('/resource/workers') export const managerService = { - getMonitors + getMonitors, + getWorkerResources } diff --git a/seatunnel-engine/seatunnel-engine-ui/src/service/manager/types.ts b/seatunnel-engine/seatunnel-engine-ui/src/service/manager/types.ts index 62bd70b910..dd8833aee5 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/service/manager/types.ts +++ b/seatunnel-engine/seatunnel-engine-ui/src/service/manager/types.ts @@ -65,3 +65,26 @@ export interface Monitor { 'client.connection.count': string 'connection.count': string } + +export interface WorkerResource { + address: string + tags?: Record<string, string> + totalSlots?: number + freeSlots?: number + usedSlots?: number + dynamicSlot?: boolean + totalCpuCores?: number | null + availableCpuCores?: number | null + totalHeapMemoryBytes?: number | null + availableHeapMemoryBytes?: number | null + cpuUsage?: number | null + memUsage?: number | null + // Only the count is displayed: JSON numbers cannot represent every Java long ID. + runningJobIds?: unknown[] +} + +export interface WorkerResourceSnapshot { + available: boolean + collectedAt: number + workers: WorkerResource[] +} diff --git a/seatunnel-engine/seatunnel-engine-ui/src/tests/managers.spec.ts b/seatunnel-engine/seatunnel-engine-ui/src/tests/managers.spec.ts index 41bc9a2229..af0995b95e 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/tests/managers.spec.ts +++ b/seatunnel-engine/seatunnel-engine-ui/src/tests/managers.spec.ts @@ -15,52 +15,366 @@ * limitations under the License. */ -import { describe, test, expect, vi, beforeEach } from 'vitest' +import { afterEach, beforeEach, describe, expect, test, vi } from 'vitest' import { flushPromises, mount } from '@vue/test-utils' -// import { createTestingPinia } from '@pinia/testing' -import { createApp } from 'vue' -import { createPinia, setActivePinia } from 'pinia' +import { createMemoryHistory, createRouter } from 'vue-router' +import { NButton, NDataTable, NDrawer } from 'naive-ui' import i18n from '@/locales' -import type { Monitor } from '@/service/manager/types' +import type { Monitor, WorkerResource, WorkerResourceSnapshot } from '@/service/manager/types' import { managerService } from '@/service/manager' import managers from '@/views/managers' +import { + addressKey, + bytesValue, + joinResources, + numberValue, + ratioValue +} from '@/views/managers/resources' -describe('managers', () => { - const app = createApp({}) - beforeEach(() => { - const pinia = createPinia() - app.use(pinia) - setActivePinia(createPinia()) - }) - test('managers component', async () => { - const mockData = [ - { - isMaster: 'true', - host: 'localhost', - port: '5801', - 'physical.memory.total': '3.6G', - 'heap.memory.used': '229.6M' - }, - { - isMaster: 'false', - host: 'localhost', - port: '5802', - 'physical.memory.total': '3.6G', - 'heap.memory.used': '1002.6M' - } - ] as Monitor[] +const monitor = (host = 'localhost', port = '5802', master = false) => + ({ + isMaster: String(master), + host, + port, + 'physical.memory.total': '3.6G', + 'heap.memory.used': '1002.6M', + 'heap.memory.max': '2G', + 'load.process': '10%', + 'thread.count': '24', + 'minor.gc.count': '2', + 'major.gc.count': '0' + }) as Monitor +const worker = (address = 'localhost:5802', dynamicSlot = false): WorkerResource => ({ + address, + dynamicSlot, + totalSlots: 4, + usedSlots: 3, + freeSlots: 1, + cpuUsage: 0, + memUsage: 0.5, + availableCpuCores: 0, + totalCpuCores: 8, + totalHeapMemoryBytes: 1048576, + availableHeapMemoryBytes: null, + runningJobIds: JSON.parse('[9223372036854775807]'), + tags: { region: 'west' } +}) +const snapshot = (workers = [worker()], available = true): WorkerResourceSnapshot => ({ + available, + collectedAt: 1723017600000, + workers +}) +const deferred = <T>() => { + let resolve!: (value: T) => void + const promise = new Promise<T>((done) => { + resolve = done + }) + return { promise, resolve } +} - vi.spyOn(managerService, 'getMonitors').mockResolvedValue(mockData) +describe('manager resource helpers', () => { + test.each([ + ['10.0.0.8', '[10.0.0.8]:5801', '10.0.0.8:5801'], + ['Worker.EXAMPLE', '[WORKER.example]:5801', 'worker.example:5801'] + ])('joins Hazelcast bracketed addresses for %s', (host, address, key) => { + const monitoring = monitor(host, '5801') + const resource = worker(address) + expect(addressKey(address)).toBe(key) + expect(joinResources([monitoring], [resource], false)).toEqual([ + { address: key, monitor: monitoring, resource } + ]) + }) + test('joins bracketed, expanded and compressed IPv6 with host and port', () => { + expect(addressKey('2001:0db8:0:0:0:0:0:1:5801')).toBe('[2001:db8::1]:5801') + expect( + joinResources( + [monitor('2001:0db8:0:0:0:0:0:1', '5801')], + [worker('[2001:db8::1]:5801')], + false + ) + ).toHaveLength(1) + expect(addressKey('[fe80::1%en0]:5801')).toBe('[fe80::1%en0]:5801') + expect(addressKey('[::1]:80')).toBe('[::1]:80') + }) + test('preserves unmatched, mixed-role and resource-only workers without inventing resources', () => { + const rows = joinResources( + [monitor('missing'), monitor('both', '5802', true)], + [worker('both:5802'), worker('resources-only:5802')], + false + ) + expect(rows).toHaveLength(3) + expect(rows.find((row) => row.address === 'missing:5802')?.resource).toBeUndefined() + expect(rows.find((row) => row.address === 'both:5802')?.monitor).toBeDefined() + expect(rows.find((row) => row.address === 'resources-only:5802')?.monitor).toBeUndefined() + expect(joinResources([monitor(), monitor('master', '5801', true)], [worker()], true)).toEqual([ + { address: 'master:5801', monitor: monitor('master', '5801', true) } + ]) + }) + test('ignores malformed records and distinguishes missing numbers from zero', () => { + expect( + joinResources( + [null, {}, monitor()] as Monitor[], + [null, {}, worker()] as WorkerResource[], + false + ) + ).toHaveLength(1) + expect(numberValue(null)).toBe('—') + expect(numberValue(0)).toBe('0') + expect(numberValue(NaN)).toBe('—') + expect(ratioValue(0)).toBe('0.0%') + expect(ratioValue(0.42)).toBe('42.0%') + expect(ratioValue(-1)).toBe('—') + expect(ratioValue(1.5)).toBe('—') + expect(bytesValue(null)).toBe('—') + expect(bytesValue(1048576)).toBe('1.0 MiB') + }) +}) - const wrapper = mount(managers, { - global: { - // plugins: [createTestingPinia({ createSpy: vi.fn() }), i18n] - plugins: [i18n] - } +describe('managers', () => { + const wrappers: ReturnType<typeof mount>[] = [] + let visibility: DocumentVisibilityState + function setVisibility(state: DocumentVisibilityState) { + visibility = state + document.dispatchEvent(new Event('visibilitychange')) + } + beforeEach(() => { + vi.useFakeTimers() + visibility = 'visible' + vi.spyOn(document, 'visibilityState', 'get').mockImplementation(() => visibility) + vi.spyOn(managerService, 'getMonitors').mockResolvedValue([monitor()]) + vi.spyOn(managerService, 'getWorkerResources').mockResolvedValue(snapshot()) + i18n.global.locale.value = 'en_US' + }) + afterEach(() => { + wrappers.forEach((wrapper) => wrapper.unmount()) + wrappers.length = 0 + vi.restoreAllMocks() + vi.useRealTimers() + document.body.innerHTML = '' + }) + async function setup(path = '/managers/workers') { + const router = createRouter({ + history: createMemoryHistory(), + routes: [{ path: '/managers/:role', component: managers }] }) - expect(managerService.getMonitors).toHaveBeenCalledTimes(1) - expect(managerService.getMonitors).toHaveBeenCalledWith() + await router.push(path) + await router.isReady() + const wrapper = mount(managers, { global: { plugins: [i18n, router] } }) + wrappers.push(wrapper) + return { wrapper, router } + } + test('renders fixed slots, monitoring and resource details without lossy job IDs', async () => { + const { wrapper } = await setup() + await flushPromises() + expect(wrapper.text()).toContain('localhost:5802') + expect(wrapper.text()).toContain('3 / 4 used, 1 free') + expect(wrapper.text()).toContain('10%') + expect(wrapper.text()).toContain('Resource response time') + await wrapper + .findAllComponents(NButton) + .find((button) => button.text() === 'Details')! + .trigger('click') + await flushPromises() + expect(document.body.textContent).toContain('0 / 8') + expect(document.body.textContent).toContain('0.0%') + expect(document.body.textContent).toContain('heap.memory.max') + expect(document.body.textContent).not.toContain('922337203685') + expect(wrapper.findComponent(NDrawer).props('show')).toBe(true) + }) + test('dynamic slots never show tracked totals as capacity', async () => { + vi.mocked(managerService.getWorkerResources).mockResolvedValue( + snapshot([worker('localhost:5802', true)]) + ) + const { wrapper } = await setup() await flushPromises() + expect(wrapper.text()).toContain('Dynamic — 3 used (no fixed capacity)') + expect(wrapper.text()).not.toContain('3 / 4') + }) + test('master page never reads worker resources or renders slots', async () => { + vi.mocked(managerService.getMonitors).mockResolvedValue([ + monitor(), + monitor('master', '5801', true) + ]) + const { wrapper } = await setup('/managers/master') + await flushPromises() + expect(wrapper.text()).toContain('master:5801') + expect(wrapper.text()).not.toContain('localhost') + expect(wrapper.text()).not.toContain('Slots') + expect(managerService.getWorkerResources).not.toHaveBeenCalled() + expect(wrapper.findComponent(NDataTable).props('columns')).toEqual( + expect.arrayContaining([expect.objectContaining({ key: 'details', fixed: 'right' })]) + ) + }) + test('worker table scrolls its action column with bounded slot summaries', async () => { + const { wrapper } = await setup() + await flushPromises() + const table = wrapper.findComponent(NDataTable) + expect(table.props('tableLayout')).toBe('fixed') + expect(table.props('scrollX')).toBe(1200) + expect(table.props('columns')).toEqual( + expect.arrayContaining([ + expect.objectContaining({ key: 'details', fixed: undefined }), + expect.objectContaining({ key: 'slots', width: 240 }) + ]) + ) + }) + test('clears failed resource snapshots and recovers on the next refresh', async () => { + const { wrapper } = await setup() + await flushPromises() + vi.mocked(managerService.getWorkerResources).mockResolvedValueOnce(snapshot([], false)) + await vi.advanceTimersByTimeAsync(30_000) + expect(wrapper.text()).toContain('Worker resources are unavailable') + expect(wrapper.text()).not.toContain('3 / 4') + expect(wrapper.text()).not.toContain('Resource response time:') expect(wrapper.text()).toContain('localhost') + await vi.advanceTimersByTimeAsync(30_000) + expect(wrapper.text()).toContain('3 / 4') + expect(wrapper.text()).not.toContain('Worker resources are unavailable') + }) + test('keeps resources when monitoring fails and handles both endpoint failures without leaking errors', async () => { + vi.mocked(managerService.getMonitors).mockRejectedValueOnce(new Error('secret server response')) + const { wrapper } = await setup() + await flushPromises() + expect(wrapper.text()).toContain('System monitoring is unavailable') + expect(wrapper.text()).toContain('3 / 4') + expect(wrapper.text()).not.toContain('secret') + vi.mocked(managerService.getMonitors).mockRejectedValueOnce(new Error('timeout')) + vi.mocked(managerService.getWorkerResources).mockRejectedValueOnce(new Error('timeout')) + await vi.advanceTimersByTimeAsync(30_000) + expect(wrapper.text()).toContain('Worker resources are unavailable') + expect(wrapper.findComponent(NDataTable).props('data')).toEqual([]) + expect(wrapper.findComponent(NDataTable).props('loading')).toBe(false) + }) + test('does not overlap slow requests, ignores obsolete routes, and fetches the current route', async () => { + const pending = deferred<Monitor[]>() + vi.mocked(managerService.getMonitors).mockReturnValueOnce(pending.promise) + const { wrapper, router } = await setup() + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + await router.push('/managers/master') + await flushPromises() + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + vi.mocked(managerService.getMonitors).mockResolvedValue([monitor('new-master', '5801', true)]) + pending.resolve([monitor('old-worker')]) + await flushPromises() + expect(wrapper.text()).toContain('new-master') + expect(wrapper.text()).not.toContain('old-worker') + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + expect(managerService.getWorkerResources).toHaveBeenCalledTimes(1) + }) + test.each(['workers', 'master'])('pauses hidden %s polling and resumes', async (role) => { + await setup(`/managers/${role}`) + await flushPromises() + await vi.advanceTimersByTimeAsync(10_000) + setVisibility('hidden') + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + setVisibility('visible') + await flushPromises() + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + await vi.advanceTimersByTimeAsync(29_999) + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + await vi.advanceTimersByTimeAsync(1) + expect(managerService.getMonitors).toHaveBeenCalledTimes(3) + expect(managerService.getWorkerResources).toHaveBeenCalledTimes(role === 'workers' ? 3 : 0) + }) + test('defers initial loading and hidden route changes until visible', async () => { + setVisibility('hidden') + const { wrapper, router } = await setup() + await router.push('/managers/master') + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).not.toHaveBeenCalled() + vi.mocked(managerService.getMonitors).mockResolvedValue([monitor('master', '5801', true)]) + setVisibility('visible') + await flushPromises() + expect(wrapper.text()).toContain('master:5801') + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + expect(managerService.getWorkerResources).not.toHaveBeenCalled() + }) + test('does not rearm polling when an in-flight request completes while hidden', async () => { + const pending = deferred<Monitor[]>() + vi.mocked(managerService.getMonitors).mockReturnValueOnce(pending.promise) + const { wrapper, router } = await setup() + setVisibility('hidden') + await router.push('/managers/master') + pending.resolve([monitor('old-worker')]) + await flushPromises() + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + expect(wrapper.findComponent(NDataTable).props('data')).toEqual([]) + vi.mocked(managerService.getMonitors).mockResolvedValue([monitor('master', '5801', true)]) + setVisibility('visible') + await flushPromises() + expect(wrapper.text()).toContain('master:5801') + expect(wrapper.text()).not.toContain('old-worker') + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + expect(managerService.getWorkerResources).toHaveBeenCalledTimes(1) + }) + test('queues one visible refresh behind an in-flight request without overlapping', async () => { + const pending = deferred<Monitor[]>() + const resumed = deferred<Monitor[]>() + vi.mocked(managerService.getMonitors) + .mockReturnValueOnce(pending.promise) + .mockReturnValueOnce(resumed.promise) + const { wrapper } = await setup() + setVisibility('hidden') + setVisibility('visible') + setVisibility('hidden') + setVisibility('visible') + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + pending.resolve([monitor('old-worker')]) + await flushPromises() + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + expect(wrapper.text()).not.toContain('old-worker') + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + resumed.resolve([monitor()]) + await flushPromises() + await vi.advanceTimersByTimeAsync(30_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(3) + expect(managerService.getWorkerResources).toHaveBeenCalledTimes(3) + }) + test('unmount removes the visibility listener and prevents restarting', async () => { + const addListener = vi.spyOn(document, 'addEventListener') + const removeListener = vi.spyOn(document, 'removeEventListener') + const { wrapper } = await setup() + await flushPromises() + const listener = addListener.mock.calls.find(([event]) => event === 'visibilitychange')?.[1] + expect(listener).toBeTypeOf('function') + setVisibility('hidden') + wrapper.unmount() + wrappers.length = 0 + expect(removeListener).toHaveBeenCalledWith('visibilitychange', listener) + setVisibility('visible') + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + expect(managerService.getWorkerResources).toHaveBeenCalledTimes(1) + }) + test('unmount ignores in-flight completion and never schedules another poll', async () => { + const pending = deferred<Monitor[]>() + vi.mocked(managerService.getMonitors).mockReturnValueOnce(pending.promise) + const { wrapper } = await setup() + wrapper.unmount() + wrappers.length = 0 + pending.resolve([monitor()]) + await flushPromises() + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(1) + }) + test('unmount cancels a scheduled poll and malformed records do not stop refresh', async () => { + vi.mocked(managerService.getMonitors).mockResolvedValueOnce([null, {}, monitor()] as Monitor[]) + vi.mocked(managerService.getWorkerResources).mockResolvedValueOnce( + snapshot([null, {}, worker()] as WorkerResource[]) + ) + const { wrapper } = await setup() + await flushPromises() + expect(wrapper.text()).toContain('3 / 4') + await vi.advanceTimersByTimeAsync(30_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) + wrapper.unmount() + wrappers.length = 0 + await vi.advanceTimersByTimeAsync(90_000) + expect(managerService.getMonitors).toHaveBeenCalledTimes(2) }) }) diff --git a/seatunnel-engine/seatunnel-engine-ui/src/views/managers/index.tsx b/seatunnel-engine/seatunnel-engine-ui/src/views/managers/index.tsx index 758a803434..4fb16109c4 100644 --- a/seatunnel-engine/seatunnel-engine-ui/src/views/managers/index.tsx +++ b/seatunnel-engine/seatunnel-engine-ui/src/views/managers/index.tsx @@ -14,80 +14,276 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -import { defineComponent, getCurrentInstance, h, ref } from 'vue' -import { useMessage, NDataTable } from 'naive-ui' -import { useI18n } from 'vue-i18n' + +import { computed, defineComponent, onBeforeUnmount, ref, watch } from 'vue' +import { + NAlert, + NButton, + NDataTable, + NDescriptions, + NDescriptionsItem, + NDrawer, + NDrawerContent, + NLayout, + NLayoutContent, + NSpace +} from 'naive-ui' import type { DataTableColumns } from 'naive-ui' -import { NButton } from 'naive-ui' -import { NSpace, NLayout, NLayoutContent } from 'naive-ui' -import { managerService } from '@/service/manager' -import type { Monitor } from '@/service/manager/types' +import { useI18n } from 'vue-i18n' import { useRoute } from 'vue-router' +import { managerService } from '@/service/manager' +import type { Monitor, WorkerResource, WorkerResourceSnapshot } from '@/service/manager/types' +import { bytesValue, joinResources, numberValue, ratioValue } from './resources' +import type { NodeResources } from './resources' export default defineComponent({ setup() { const { t } = useI18n() const route = useRoute() - const monitors = ref([] as Monitor[]) + const isMaster = computed(() => route?.path.endsWith('/master') || false) + const rows = ref<NodeResources[]>([]) + const loading = ref(false) + const monitorUnavailable = ref(false) + const resourceUnavailable = ref(false) + const collectedAt = ref<number>() + const selectedAddress = ref<string>() + const selected = computed(() => rows.value.find((row) => row.address === selectedAddress.value)) + let timer: ReturnType<typeof setTimeout> | undefined + let generation = 0 + let disposed = false + const isHidden = () => document.visibilityState === 'hidden' + + const refresh = async () => { + if (disposed || isHidden() || loading.value) return + clearTimeout(timer) + loading.value = true + const requestGeneration = generation + const master = isMaster.value + try { + const [monitorResult, resourceResult] = await Promise.allSettled([ + managerService.getMonitors(), + master ? Promise.resolve(undefined) : managerService.getWorkerResources() + ]) + if (!disposed && requestGeneration === generation) { + const monitors: Monitor[] = + monitorResult.status === 'fulfilled' && Array.isArray(monitorResult.value) + ? monitorResult.value + : [] + monitorUnavailable.value = + monitorResult.status === 'rejected' || !Array.isArray(monitorResult.value) + const snapshot: WorkerResourceSnapshot | undefined = + resourceResult.status === 'fulfilled' ? resourceResult.value : undefined + const available = snapshot?.available === true && Array.isArray(snapshot.workers) + resourceUnavailable.value = !master && !available + // Replace, rather than retain, old resource values after failed refreshes. + rows.value = joinResources(monitors, available ? snapshot!.workers : [], master) + collectedAt.value = + available && Number.isFinite(snapshot!.collectedAt) && snapshot!.collectedAt > 0 + ? snapshot!.collectedAt + : undefined + } + } finally { + loading.value = false + if (!disposed && !isHidden()) { + if (requestGeneration !== generation) void refresh() + else timer = setTimeout(refresh, 30_000) + } + } + } - const fetch = async () => { - let res = await managerService.getMonitors() - const isMaster = route?.path.endsWith('/master') || false - res = res.filter((row) => row.isMaster === String(isMaster)) || [] - monitors.value = res + const onVisibilityChange = () => { + clearTimeout(timer) + if (isHidden()) { + // Discard the pending sample and refresh after it settles if visibility returns first. + generation++ + } else { + void refresh() + } } - fetch() + document.addEventListener('visibilitychange', onVisibilityChange) - function createColumns(): DataTableColumns<Monitor> { - const view = (row: Monitor) => {} - return [ - { - title: 'Host', - key: 'host' - }, - { - title: 'Port', - key: 'port' - }, - { - title: 'Physical MEM', - key: 'physical.memory.total' - }, - { - title: 'Heap MEM Used', - key: 'heap.memory.used' - // }, - // { - // title: 'Action', - // key: 'actions', - // render(row) { - // return h( - // NButton, - // { - // strong: true, - // tertiary: true, - // size: 'small', - // onClick: () => view(row) - // }, - // { default: () => 'View' } - // ) - // } - } - ] + watch( + isMaster, + () => { + generation++ + clearTimeout(timer) + rows.value = [] + selectedAddress.value = undefined + collectedAt.value = undefined + monitorUnavailable.value = false + resourceUnavailable.value = false + void refresh() + }, + { immediate: true } + ) + onBeforeUnmount(() => { + disposed = true + clearTimeout(timer) + document.removeEventListener('visibilitychange', onVisibilityChange) + }) + + const monitorValue = (row: NodeResources, field: keyof Monitor) => row.monitor?.[field] ?? '—' + const slots = (resource?: WorkerResource) => { + if (typeof resource?.dynamicSlot !== 'boolean') return '—' + if (resource.dynamicSlot) { + return t('managers.dynamic_used', { used: numberValue(resource.usedSlots) }) + } + return t('managers.fixed_slots', { + used: numberValue(resource.usedSlots), + total: numberValue(resource.totalSlots), + free: numberValue(resource.freeSlots) + }) } + const columns = computed<DataTableColumns<NodeResources>>(() => [ + { title: t('managers.address'), key: 'address' }, + { title: t('managers.cpu'), key: 'cpu', render: (row) => monitorValue(row, 'load.process') }, + { + title: t('managers.heap'), + key: 'heap', + render: (row) => + `${monitorValue(row, 'heap.memory.used')} / ${monitorValue(row, 'heap.memory.max')}` + }, + { + title: t('managers.physical'), + key: 'physical', + render: (row) => monitorValue(row, 'physical.memory.total') + }, + { + title: t('managers.gc'), + key: 'gc', + render: (row) => + `${monitorValue(row, 'minor.gc.count')} / ${monitorValue(row, 'major.gc.count')}` + }, + { + title: t('managers.threads'), + key: 'threads', + render: (row) => monitorValue(row, 'thread.count') + }, + ...(!isMaster.value + ? [ + { + title: t('managers.slots'), + key: 'slots', + width: 240, + render: (row: NodeResources) => ( + <span style={{ display: 'block', whiteSpace: 'normal', overflowWrap: 'anywhere' }}> + {slots(row.resource)} + </span> + ) + } + ] + : []), + { + title: t('managers.details'), + key: 'details', + fixed: isMaster.value ? 'right' : undefined, + width: 95, + render: (row) => ( + <NButton + size="small" + onClick={() => { + selectedAddress.value = row.address + }} + > + {t('managers.details')} + </NButton> + ) + } + ]) + const resourceDetails = (resource: WorkerResource) => [ + [t('managers.slots'), slots(resource)], + [ + t('managers.cpu_resources'), + `${numberValue(resource.availableCpuCores)} / ${numberValue(resource.totalCpuCores)}` + ], + [ + t('managers.heap_resources'), + `${bytesValue(resource.availableHeapMemoryBytes)} / ${bytesValue(resource.totalHeapMemoryBytes)}` + ], + [t('managers.cpu_usage'), ratioValue(resource.cpuUsage)], + [t('managers.memory_usage'), ratioValue(resource.memUsage)], + [ + t('managers.running_jobs'), + Array.isArray(resource.runningJobIds) ? String(resource.runningJobIds.length) : '—' + ], + [t('managers.tags'), resource.tags ? JSON.stringify(resource.tags) : '—'] + ] - const columns = createColumns() return () => ( <NLayout> <NLayoutContent> <div class="w-full bg-white p-6 border border-gray-100 rounded-xl"> - <h2 class="font-bold text-2xl pb-6">{t('managers.managers')}</h2> + <NSpace justify="space-between"> + <h2 class="font-bold text-2xl pb-6">{t('managers.managers')}</h2> + <NButton + loading={loading.value} + disabled={loading.value} + onClick={() => void refresh()} + > + {t('managers.refresh')} + </NButton> + </NSpace> + <p class="pb-3">{t('managers.refresh_hint')}</p> + {monitorUnavailable.value && ( + <NAlert type="warning">{t('managers.monitor_unavailable')}</NAlert> + )} + {resourceUnavailable.value && ( + <NAlert type="warning">{t('managers.resource_unavailable')}</NAlert> + )} + {!isMaster.value && ( + <p class="py-3"> + {t('managers.snapshot_hint')} + {collectedAt.value && + ` ${t('managers.collected_at')}: ${new Date(collectedAt.value).toLocaleString()}`} + </p> + )} <NDataTable - columns={columns} - data={monitors.value} - pagination={false} + columns={columns.value} + data={rows.value} + loading={loading.value} + rowKey={(row: NodeResources) => row.address} + pagination={{ pageSize: 20 }} + tableLayout={isMaster.value ? 'auto' : 'fixed'} + scrollX={1200} bordered={false} /> + <NDrawer + show={!!selected.value} + width="min(640px, 100vw)" + onUpdateShow={(show) => { + if (!show) selectedAddress.value = undefined + }} + > + <NDrawerContent title={selected.value?.address} closable> + {selected.value?.resource && ( + <> + <h3>{t('managers.resources')}</h3> + <NDescriptions column={1} bordered> + {resourceDetails(selected.value.resource).map(([label, value]) => ( + <NDescriptionsItem key={label} label={label}> + {value} + </NDescriptionsItem> + ))} + </NDescriptions> + </> + )} + {!isMaster.value && !selected.value?.resource && ( + <NAlert type="warning">{t('managers.resource_missing')}</NAlert> + )} + <h3 class="py-3">{t('managers.monitoring')}</h3> + {selected.value?.monitor ? ( + <NDescriptions column={1} bordered> + {Object.entries(selected.value.monitor).map(([field, value]) => ( + <NDescriptionsItem key={field} label={field}> + {String(value ?? '—')} + </NDescriptionsItem> + ))} + </NDescriptions> + ) : ( + <NAlert type="warning">{t('managers.monitor_missing')}</NAlert> + )} + </NDrawerContent> + </NDrawer> </div> </NLayoutContent> </NLayout> diff --git a/seatunnel-engine/seatunnel-engine-ui/src/views/managers/resources.ts b/seatunnel-engine/seatunnel-engine-ui/src/views/managers/resources.ts new file mode 100644 index 0000000000..26873aa969 --- /dev/null +++ b/seatunnel-engine/seatunnel-engine-ui/src/views/managers/resources.ts @@ -0,0 +1,85 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +import type { Monitor, WorkerResource } from '@/service/manager/types' + +export interface NodeResources { + address: string + monitor?: Monitor + resource?: WorkerResource +} + +// Hazelcast brackets IPv6 addresses; monitoring returns host and port separately. +export function addressKey(address: string): string { + const separator = address.lastIndexOf(':') + if (separator < 0) return address + const host = address.slice(0, separator).replace(/^\[|\]$/g, '') + const port = address.slice(separator + 1) + if (host.includes(':')) { + try { + return `${new URL(`http://[${host}]:${port}`).hostname}:${port}` + } catch { + return `[${host}]:${port}` + } + } + return `${host.toLowerCase()}:${port}` +} + +export function joinResources( + monitors: Monitor[], + workers: WorkerResource[], + master: boolean +): NodeResources[] { + monitors = monitors.filter( + (monitor) => + monitor && + typeof monitor.host === 'string' && + (typeof monitor.port === 'string' || typeof monitor.port === 'number') + ) + const rows = new Map<string, NodeResources>() + for (const monitor of monitors) { + const address = addressKey(`${monitor.host}:${monitor.port}`) + if (monitor.isMaster === String(master)) rows.set(address, { address, monitor }) + } + if (!master) { + const byAddress = new Map( + monitors.map((monitor) => [addressKey(`${monitor.host}:${monitor.port}`), monitor]) + ) + for (const resource of workers) { + if (!resource || typeof resource.address !== 'string') continue + const address = addressKey(resource.address) + // A mixed-role member can be both the master and a registered worker. + rows.set(address, { address, monitor: byAddress.get(address), resource }) + } + } + return Array.from(rows.values()) +} + +export function numberValue(value: number | null | undefined): string { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 ? String(value) : '—' +} + +export function ratioValue(value: number | null | undefined): string { + return typeof value === 'number' && Number.isFinite(value) && value >= 0 && value <= 1 + ? `${(value * 100).toFixed(1)}%` + : '—' +} + +export function bytesValue(value: number | null | undefined): string { + if (typeof value !== 'number' || !Number.isFinite(value) || value < 0) return '—' + return `${(value / 1024 / 1024).toFixed(1)} MiB` +}
