This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-8552-e321f0c717cf5873cebd3ce72885f41976b0282d in repository https://gitbox.apache.org/repos/asf/texera.git
commit 51d7564309bf74ad75b4f28844ba2516e6bad9a6 Author: Prateek Ganigi <[email protected]> AuthorDate: Mon Sep 21 01:21:07 2026 +0000 feat(frontend): restore the heat-map after a page refresh (#8552) ### What changes were proposed in this PR? Third sub-task of the heat-map umbrella #5772. Today the heat-map is session-only: a refresh drops the Layers > Performance toggle and the statistics that color the canvas, so the workflow has to be re-run just to look at it again. **Before this PR:** https://github.com/user-attachments/assets/3a29a5b1-5e7f-4d7a-9e62-05efdb47d6f1 **After this PR:** https://github.com/user-attachments/assets/a2506966-dbd5-4c93-8304-cb950d9a2771 ``` Before: refresh -> session state cleared -> overlay off, canvas blank After: refresh -> toggle + view restored -> last run's stats re-fetched -> canvas repainted ``` Two halves, both frontend-only. **1. Overlay state.** `heatmap-overlay-persistence.ts` keeps `{ view: HeatmapView | null }` in localStorage via the existing `localSetObject` / `localGetObject`. `MenuComponent` saves on toggle and on view change, and restores in `ngOnInit`. Only the Performance layer persists: Grid / Regions / Workers / Status stay session-only. A corrupted or foreign value reads as off. **2. Statistics.** `HeatmapStatsRestoreService` composes existing client APIs: ``` wid -> retrieveWorkflowExecutions -> latest completed run -> retrieveWorkflowRuntimeStatistics -> latest snapshot per operator -> WorkflowStatusService.setExternalStatus ``` `runtime-statistics-mapper.ts` does the shape conversion, including status code -> `OperatorState`; Failed and Killed have no `OperatorState` member, so they fall back to `Uninitialized` rather than a wrong state. Port-level metrics are not persisted, so the mapper omits the port maps entirely rather than emitting empty ones: `{}` already means "every port measured zero" (what `resetStatus` emits), so `changeOperatorStatistics` skips the port loops when the maps are absent and the port labels keep their display names. Two changes land outside the heat-map feature. `inputPortMetrics` / `outputPortMetrics` become optional on `OperatorStatistics`, the shared websocket-mirror type. Jackson already omits absent fields, so the mirror is not weakened. Separately, `changeOperatorStatistics` gains a branch for the absent case. `WorkflowStatusService` gains `setExternalStatus`. The websocket handler's body moves into a private `ingestRuntimeStatus` that both paths share, so live and restored updates emit identically. `performanceMetricsSubject` is already a `BehaviorSubject`, so the restore fetch cannot race the overlay's own init. Restoring is best-effort and entirely gated inside the service: | Skips when | Why | | --- | --- | | Overlay not persisted on | A localStorage read, so non-users issue no request at all | | No `wid` (unsaved workflow) | Nothing to fetch | | Execution in progress | The live stream wins | | No executions, or no statistics | Nothing to restore | | HTTP error | Must never block workspace entry | | Another producer writes statistics first | A new run resets the canvas; `isExecuting()` cannot see it until the backend replies | ### Any related issues, documentation, discussions? Closes #5775. Part of umbrella #5772. RFC #5216. Reviewer note: #5775 refers to `setExternalStatus` as already existing from #5773. It is not on `main` and this PR adds it. ### How was this PR tested? Manual, Chrome: run a workflow to completion, Layers > Performance on, pick **I/O imbalance**, reload: toggle, view and colors all come back with no re-run. Switch to **Time / row** and reload: the ranking inverts and survives. Uncheck and reload: stays off. `localStorage.clear()` and reload: back to default. ``` cd frontend && ng test --watch=false \ --include "**/heatmap-overlay-persistence.spec.ts" \ --include "**/runtime-statistics-mapper.spec.ts" \ --include "**/heatmap-stats-restore.service.spec.ts" \ --include "**/workflow-status.service.spec.ts" \ --include "**/joint-ui.service.spec.ts" \ --include "**/menu.component.spec.ts" \ --include "**/workspace.component.spec.ts" ``` | Spec | Tests | New | | --- | --- | --- | | `heatmap-overlay-persistence.spec.ts` | 4 | 4 | | `runtime-statistics-mapper.spec.ts` | 8 | 8 | | `heatmap-stats-restore.service.spec.ts` | 16 | 16 | | `workflow-status.service.spec.ts` | 14 | 3 | | `joint-ui.service.spec.ts` | 71 | 2 | | `menu.component.spec.ts` | 128 | 6 | | `workspace.component.spec.ts` | 28 | 0 | 39 new tests, 269 passing, covering the negative paths: corrupted stored value, unsaved workflow, live execution, a run started while the fetches are in flight, zero executions, a run that never completed, an empty payload, and an HTTP error from either fetch. `tsc --noEmit` (strict), `eslint ./src` and Prettier are clean. ### Was this PR authored or co-authored using generative AI tooling? This PR was co-authored by Claude, in compliance with ASF. --- .../component/menu/menu.component.spec.ts | 66 +++++ .../app/workspace/component/menu/menu.component.ts | 23 +- .../app/workspace/component/workspace.component.ts | 15 +- .../heatmap/heatmap-overlay-persistence.spec.ts | 58 ++++ .../service/heatmap/heatmap-overlay-persistence.ts | 51 ++++ .../heatmap/heatmap-stats-restore.service.spec.ts | 295 +++++++++++++++++++++ .../heatmap/heatmap-stats-restore.service.ts | 124 +++++++++ .../heatmap/runtime-statistics-mapper.spec.ts | 130 +++++++++ .../service/heatmap/runtime-statistics-mapper.ts | 78 ++++++ .../service/joint-ui/joint-ui.service.spec.ts | 35 +++ .../workspace/service/joint-ui/joint-ui.service.ts | 42 +-- .../workflow-status.service.spec.ts | 36 +++ .../workflow-status/workflow-status.service.ts | 35 ++- .../workspace/types/execute-workflow.interface.ts | 5 +- 14 files changed, 961 insertions(+), 32 deletions(-) diff --git a/frontend/src/app/workspace/component/menu/menu.component.spec.ts b/frontend/src/app/workspace/component/menu/menu.component.spec.ts index 1214f513af..0ed34ae0ec 100644 --- a/frontend/src/app/workspace/component/menu/menu.component.spec.ts +++ b/frontend/src/app/workspace/component/menu/menu.component.spec.ts @@ -43,6 +43,7 @@ import { WorkflowPersistService } from "../../../common/service/workflow-persist import { NotificationService } from "../../../common/service/notification/notification.service"; import { ExecutionState } from "../../types/execute-workflow.interface"; import { HeatmapView } from "../../service/heatmap/heatmap-scoring"; +import { loadPersistedHeatmapView, savePersistedHeatmapView } from "../../service/heatmap/heatmap-overlay-persistence"; import { ComputingUnitState } from "../../../common/type/computing-unit-connection.interface"; import { mockPoint, mockScanPredicate } from "../../service/workflow-graph/model/mock-workflow-data"; import { FileSaverService } from "../../../dashboard/service/user/file/file-saver.service"; @@ -1088,6 +1089,71 @@ describe("MenuComponent", () => { }); }); + describe("heat-map overlay persistence", () => { + beforeEach(() => localStorage.clear()); + afterEach(() => localStorage.clear()); + + it("restores a persisted view on init: checkbox on, selector set, view pushed", () => { + savePersistedHeatmapView(HeatmapView.IoImbalance); + const setSpy = vi.spyOn(workflowActionService.getJointGraphWrapper(), "setHeatmapView"); + + component.restorePersistedHeatmapOverlay(); + + expect(component.showHeatmap).toBe(true); + expect(component.heatmapView).toBe(HeatmapView.IoImbalance); + expect(setSpy).toHaveBeenCalledWith(HeatmapView.IoImbalance); + }); + + it("leaves the overlay off when nothing is persisted", () => { + const setSpy = vi.spyOn(workflowActionService.getJointGraphWrapper(), "setHeatmapView"); + + component.restorePersistedHeatmapOverlay(); + + expect(component.showHeatmap).toBe(false); + expect(setSpy).not.toHaveBeenCalled(); + }); + + it("leaves the overlay off for a persisted null view", () => { + savePersistedHeatmapView(null); + const setSpy = vi.spyOn(workflowActionService.getJointGraphWrapper(), "setHeatmapView"); + + component.restorePersistedHeatmapOverlay(); + + expect(component.showHeatmap).toBe(false); + expect(setSpy).not.toHaveBeenCalled(); + }); + + it("persists the view when the overlay is toggled on, and null when toggled off", () => { + component.showHeatmap = true; + component.heatmapView = HeatmapView.TimePerRow; + component.toggleHeatmap(); + expect(loadPersistedHeatmapView()).toBe(HeatmapView.TimePerRow); + + component.showHeatmap = false; + component.toggleHeatmap(); + expect(loadPersistedHeatmapView()).toBeNull(); + }); + + it("persists a view change made while the overlay is on", () => { + component.showHeatmap = true; + component.heatmapView = HeatmapView.Runtime; + component.toggleHeatmap(); + + component.setHeatmapView(HeatmapView.IoImbalance); + + expect(loadPersistedHeatmapView()).toBe(HeatmapView.IoImbalance); + }); + + it("does not persist a view browsed while the overlay is off", () => { + component.showHeatmap = false; + component.toggleHeatmap(); + + component.setHeatmapView(HeatmapView.IoImbalance); + + expect(loadPersistedHeatmapView()).toBeNull(); + }); + }); + describe("toggleStatus", () => { it("removes hide-operator-status when enabled and repositions the status label", () => { const operator = fakeElement("operator"); diff --git a/frontend/src/app/workspace/component/menu/menu.component.ts b/frontend/src/app/workspace/component/menu/menu.component.ts index cc1fa2c727..865e147ed0 100644 --- a/frontend/src/app/workspace/component/menu/menu.component.ts +++ b/frontend/src/app/workspace/component/menu/menu.component.ts @@ -29,6 +29,7 @@ import { ValidationWorkflowService } from "../../service/validation/validation-w import { WorkflowActionService } from "../../service/workflow-graph/model/workflow-action.service"; import { ExecutionState } from "../../types/execute-workflow.interface"; import { HeatmapView } from "../../service/heatmap/heatmap-scoring"; +import { loadPersistedHeatmapView, savePersistedHeatmapView } from "../../service/heatmap/heatmap-overlay-persistence"; import { WorkflowWebsocketService } from "../../service/workflow-websocket/workflow-websocket.service"; import { WorkflowResultExportService } from "../../service/workflow-result-export/workflow-result-export.service"; import { catchError, debounceTime, switchMap, tap } from "rxjs/operators"; @@ -225,6 +226,8 @@ export class MenuComponent implements OnInit, OnDestroy { } public ngOnInit(): void { + this.restorePersistedHeatmapOverlay(); + // Marks an edit for the Form View hand-over (see onClickOpenFormView): set the moment an edit is // reported, before the autosave debounce, cleared when the switch's save snapshots the workflow. this.workflowActionService @@ -602,14 +605,32 @@ export class MenuComponent implements OnInit, OnDestroy { public toggleHeatmap(): void { // The editor subscribes to this stream and colors operator fills (canvas + mini-map). // A null view turns the overlay off; a view enables it. - this.workflowActionService.getJointGraphWrapper().setHeatmapView(this.showHeatmap ? this.heatmapView : null); + const view = this.showHeatmap ? this.heatmapView : null; + this.workflowActionService.getJointGraphWrapper().setHeatmapView(view); + savePersistedHeatmapView(view); } public setHeatmapView(view: HeatmapView): void { this.heatmapView = view; if (this.showHeatmap) { this.workflowActionService.getJointGraphWrapper().setHeatmapView(view); + savePersistedHeatmapView(view); + } + } + + /** + * Restores the persisted heat-map overlay state (Layers > Performance) on + * workspace entry. Only the Performance layer persists; the other canvas + * layers stay session-only. + */ + public restorePersistedHeatmapOverlay(): void { + const view = loadPersistedHeatmapView(); + if (view === null) { + return; } + this.showHeatmap = true; + this.heatmapView = view; + this.workflowActionService.getJointGraphWrapper().setHeatmapView(view); } /** diff --git a/frontend/src/app/workspace/component/workspace.component.ts b/frontend/src/app/workspace/component/workspace.component.ts index f83f43eae5..4a26c2bc03 100644 --- a/frontend/src/app/workspace/component/workspace.component.ts +++ b/frontend/src/app/workspace/component/workspace.component.ts @@ -53,6 +53,7 @@ import { USER_WORKSPACE } from "../../app-routing.constant"; import { GuiConfigService } from "../../common/service/gui-config.service"; import { ComputingUnitStatusService } from "../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; import { ExecuteWorkflowService } from "../service/execute-workflow/execute-workflow.service"; +import { HeatmapStatsRestoreService } from "../service/heatmap/heatmap-stats-restore.service"; import { WorkflowResultService } from "../service/workflow-result/workflow-result.service"; import { checkIfWorkflowBroken } from "../../common/util/workflow-check"; import { NzSpinComponent } from "ng-zorro-antd/spin"; @@ -131,7 +132,8 @@ export class WorkspaceComponent implements AfterViewInit, OnInit, OnDestroy { private changeDetectorRef: ChangeDetectorRef, private computingUnitStatusService: ComputingUnitStatusService, private executeWorkflowService: ExecuteWorkflowService, - private workflowResultService: WorkflowResultService + private workflowResultService: WorkflowResultService, + private heatmapStatsRestoreService: HeatmapStatsRestoreService ) {} ngOnInit() { @@ -277,6 +279,17 @@ export class WorkspaceComponent implements AfterViewInit, OnInit, OnDestroy { if (shouldAutoLayout) { this.workflowActionService.autoLayoutWorkflow(); } + // Best-effort: repaint the heat-map from the last run's persisted + // statistics. The service itself gates on the overlay being + // persisted-on (and skips during a live execution), so this is a + // no-op for everyone else. Fetch failures are swallowed there; an + // error arriving here is a mapping or ingestion bug, so it is logged. + this.heatmapStatsRestoreService + .restoreLatestRunStatistics() + .pipe(untilDestroyed(this)) + .subscribe({ + error: (err: unknown) => console.error("Failed to restore heat-map statistics:", err), + }); // set the URL fragment to previous value // because reloadWorkflow will highlight/unhighlight all elements // which will change the URL fragment diff --git a/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.spec.ts b/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.spec.ts new file mode 100644 index 0000000000..04830e39ae --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.spec.ts @@ -0,0 +1,58 @@ +/** + * 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 { + HEATMAP_OVERLAY_STORAGE_KEY, + loadPersistedHeatmapView, + savePersistedHeatmapView, +} from "./heatmap-overlay-persistence"; +import { HeatmapView } from "./heatmap-scoring"; + +describe("heatmap overlay persistence", () => { + beforeEach(() => localStorage.clear()); + afterEach(() => localStorage.clear()); + + it("round-trips an active view", () => { + savePersistedHeatmapView(HeatmapView.TimePerRow); + expect(loadPersistedHeatmapView()).toBe(HeatmapView.TimePerRow); + }); + + it("round-trips the overlay-off state (null view)", () => { + savePersistedHeatmapView(HeatmapView.Runtime); + savePersistedHeatmapView(null); + expect(loadPersistedHeatmapView()).toBeNull(); + }); + + it("reads as off when nothing was ever persisted", () => { + expect(loadPersistedHeatmapView()).toBeNull(); + }); + + it("reads a corrupted or foreign stored value as off", () => { + localStorage.setItem(HEATMAP_OVERLAY_STORAGE_KEY, "not-json{"); + expect(loadPersistedHeatmapView()).toBeNull(); + + // Valid JSON, but not a HeatmapView member (e.g. persisted by a future or + // older build): must not leak an unknown value into the view stream. + localStorage.setItem(HEATMAP_OVERLAY_STORAGE_KEY, JSON.stringify({ view: "no-such-view" })); + expect(loadPersistedHeatmapView()).toBeNull(); + + localStorage.setItem(HEATMAP_OVERLAY_STORAGE_KEY, JSON.stringify({ wrongShape: true })); + expect(loadPersistedHeatmapView()).toBeNull(); + }); +}); diff --git a/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.ts b/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.ts new file mode 100644 index 0000000000..2cffe484a0 --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/heatmap-overlay-persistence.ts @@ -0,0 +1,51 @@ +/** + * 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 { localGetObject, localSetObject } from "../../../common/util/storage"; +import { HeatmapView } from "./heatmap-scoring"; + +/** localStorage key for the persisted heat-map overlay state. */ +export const HEATMAP_OVERLAY_STORAGE_KEY = "heatmap-overlay"; + +/** Persisted shape: the active view, or null when the overlay is off. */ +export interface PersistedHeatmapOverlay { + readonly view: HeatmapView | null; +} + +/** Persists the overlay state (a null view means the overlay is off). */ +export function savePersistedHeatmapView(view: HeatmapView | null): void { + localSetObject<PersistedHeatmapOverlay>(HEATMAP_OVERLAY_STORAGE_KEY, { view }); +} + +/** + * Reads the persisted overlay state back. Returns null (overlay off) when + * nothing was persisted or the stored value is not a valid HeatmapView — + * a corrupted or foreign value must not leak into the view stream. + */ +export function loadPersistedHeatmapView(): HeatmapView | null { + try { + // localGetObject JSON.parses the stored string, which throws on a + // corrupted value; treat that the same as nothing persisted. + const persisted = localGetObject<PersistedHeatmapOverlay>(HEATMAP_OVERLAY_STORAGE_KEY); + const view = persisted?.view; + return Object.values(HeatmapView).includes(view as HeatmapView) ? (view as HeatmapView) : null; + } catch { + return null; + } +} diff --git a/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.spec.ts b/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.spec.ts new file mode 100644 index 0000000000..f017a21c0e --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.spec.ts @@ -0,0 +1,295 @@ +/** + * 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 { TestBed } from "@angular/core/testing"; +import { Subject, of, throwError } from "rxjs"; +import type { Mocked } from "vitest"; +import { HeatmapStatsRestoreService } from "./heatmap-stats-restore.service"; +import { savePersistedHeatmapView } from "./heatmap-overlay-persistence"; +import { HeatmapView } from "./heatmap-scoring"; +import { WorkflowExecutionsService } from "../../../dashboard/service/user/workflow-executions/workflow-executions.service"; +import { WorkflowActionService } from "../workflow-graph/model/workflow-action.service"; +import { WorkflowStatusService } from "../workflow-status/workflow-status.service"; +import { ExecuteWorkflowService } from "../execute-workflow/execute-workflow.service"; +import { ExecutionState, OperatorState } from "../../types/execute-workflow.interface"; +import { WorkflowExecutionsEntry } from "../../../dashboard/type/workflow-executions-entry"; +import { WorkflowRuntimeStatistics } from "../../../dashboard/type/workflow-runtime-statistics"; + +function makeExecution(overrides: Partial<WorkflowExecutionsEntry>): WorkflowExecutionsEntry { + return { + eId: 1, + vId: 1, + cuId: 7, + whId: null, + sId: 0, + userName: "user", + avatar: "", + name: "run", + startingTime: 1_000, + completionTime: 2_000, + status: 3, // Completed + result: "", + bookmarked: false, + logLocation: "", + ...overrides, + }; +} + +function makeStatsRow(overrides: Partial<WorkflowRuntimeStatistics>): WorkflowRuntimeStatistics { + return { + operatorId: "op1", + timestamp: 1_000, + inputTupleCount: 10, + inputTupleSize: 100, + outputTupleCount: 5, + outputTupleSize: 50, + totalDataProcessingTime: 1_000_000, + totalControlProcessingTime: 2_000, + totalIdleTime: 3_000, + numberOfWorkers: 2, + status: 3, + ...overrides, + }; +} + +describe("HeatmapStatsRestoreService", () => { + let service: HeatmapStatsRestoreService; + let executionsService: Mocked<WorkflowExecutionsService>; + let statusService: Mocked<WorkflowStatusService>; + let actionService: Mocked<WorkflowActionService>; + let executeService: Mocked<ExecuteWorkflowService>; + let statisticsUpdates: Subject<Record<string, never>>; + + beforeEach(() => { + localStorage.clear(); + // Default arrangement: overlay persisted on, a saved workflow, no live run. + savePersistedHeatmapView(HeatmapView.Runtime); + + executionsService = { + retrieveWorkflowExecutions: vi.fn(() => of([makeExecution({})])), + retrieveWorkflowRuntimeStatistics: vi.fn(() => of([makeStatsRow({})])), + } as unknown as Mocked<WorkflowExecutionsService>; + actionService = { + getWorkflowMetadata: vi.fn(() => ({ wid: 42 })), + } as unknown as Mocked<WorkflowActionService>; + // A plain Subject, like the real one: subscribing does not emit, so it only cuts the + // restore short when another producer actually writes statistics. + statisticsUpdates = new Subject<Record<string, never>>(); + statusService = { + setExternalStatus: vi.fn(), + getStatisticsUpdateStream: vi.fn(() => statisticsUpdates.asObservable()), + } as unknown as Mocked<WorkflowStatusService>; + executeService = { + getExecutionState: vi.fn(() => ({ state: ExecutionState.Uninitialized })), + } as unknown as Mocked<ExecuteWorkflowService>; + + TestBed.configureTestingModule({ + providers: [ + HeatmapStatsRestoreService, + { provide: WorkflowExecutionsService, useValue: executionsService }, + { provide: WorkflowActionService, useValue: actionService }, + { provide: WorkflowStatusService, useValue: statusService }, + { provide: ExecuteWorkflowService, useValue: executeService }, + ], + }); + service = TestBed.inject(HeatmapStatsRestoreService); + }); + + afterEach(() => localStorage.clear()); + + it("fetches the latest run and feeds the mapped statistics into WorkflowStatusService", () => { + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).toHaveBeenCalledWith(42); + expect(executionsService.retrieveWorkflowRuntimeStatistics).toHaveBeenCalledWith(42, 1, 7); + expect(statusService.setExternalStatus).toHaveBeenCalledWith({ + op1: expect.objectContaining({ + operatorState: OperatorState.Completed, + aggregatedInputRowCount: 10, + aggregatedOutputRowCount: 5, + aggregatedDataProcessingTime: 1_000_000, + numWorkers: 2, + }), + }); + }); + + it("does nothing when the overlay is not persisted on", () => { + savePersistedHeatmapView(null); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).not.toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("does nothing for an unsaved workflow (no wid)", () => { + actionService.getWorkflowMetadata.mockReturnValue({ wid: undefined } as never); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).not.toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("skips the restore while an execution is in progress, so the live stream wins", () => { + executeService.getExecutionState.mockReturnValue({ state: ExecutionState.Running } as never); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).not.toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("is cold: reads no gate and issues no request until subscribed", () => { + const pending = service.restoreLatestRunStatistics(); + + expect(executeService.getExecutionState).not.toHaveBeenCalled(); + expect(executionsService.retrieveWorkflowExecutions).not.toHaveBeenCalled(); + + pending.subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).toHaveBeenCalledWith(42); + }); + + it("evaluates the gates at subscribe time, not at call time", () => { + const pending = service.restoreLatestRunStatistics(); + savePersistedHeatmapView(null); + + pending.subscribe(); + + expect(executionsService.retrieveWorkflowExecutions).not.toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("drops the restore when a run starts while the fetches are in flight", () => { + // The entry guard passes on Uninitialized, which is also the state before the websocket + // connects; a live run landing mid-flight must not be overwritten by the previous run. + executionsService.retrieveWorkflowRuntimeStatistics.mockImplementation(() => { + executeService.getExecutionState.mockReturnValue({ state: ExecutionState.Running } as never); + return of([makeStatsRow({})]); + }); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowRuntimeStatistics).toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("neither ingests nor throws when the workflow has no executions", () => { + executionsService.retrieveWorkflowExecutions.mockReturnValue(of([])); + + expect(() => service.restoreLatestRunStatistics().subscribe()).not.toThrow(); + + expect(executionsService.retrieveWorkflowRuntimeStatistics).not.toHaveBeenCalled(); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("prefers the most recent completed run over newer unfinished ones", () => { + executionsService.retrieveWorkflowExecutions.mockReturnValue( + of([ + makeExecution({ eId: 3, cuId: 9, startingTime: 3_000, status: 4 }), // newest, Failed + makeExecution({ eId: 2, cuId: 8, startingTime: 2_000, status: 3 }), // newest Completed + makeExecution({ eId: 1, cuId: 7, startingTime: 1_000, status: 3 }), + ]) + ); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowRuntimeStatistics).toHaveBeenCalledWith(42, 2, 8); + }); + + it("still loads the latest run when no execution ever completed", () => { + executionsService.retrieveWorkflowExecutions.mockReturnValue( + of([ + makeExecution({ eId: 5, cuId: 11, startingTime: 5_000, status: 1 }), // Running + makeExecution({ eId: 4, cuId: 10, startingTime: 4_000, status: 4 }), // Failed + ]) + ); + + service.restoreLatestRunStatistics().subscribe(); + + expect(executionsService.retrieveWorkflowRuntimeStatistics).toHaveBeenCalledWith(42, 5, 11); + }); + + // The error path is pinned via an explicit { error, complete } observer. RxJS 7 reports an + // unhandled error asynchronously instead of rethrowing from subscribe(), so a not.toThrow() + // assertion would pass with catchError deleted. + it("swallows an HTTP error from the executions fetch", () => { + executionsService.retrieveWorkflowExecutions.mockReturnValue(throwError(() => new Error("500"))); + const onError = vi.fn(); + const onComplete = vi.fn(); + + service.restoreLatestRunStatistics().subscribe({ error: onError, complete: onComplete }); + + expect(onError).not.toHaveBeenCalled(); + expect(onComplete).toHaveBeenCalledTimes(1); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("swallows an HTTP error from the statistics fetch", () => { + executionsService.retrieveWorkflowRuntimeStatistics.mockReturnValue(throwError(() => new Error("500"))); + const onError = vi.fn(); + const onComplete = vi.fn(); + + service.restoreLatestRunStatistics().subscribe({ error: onError, complete: onComplete }); + + expect(onError).not.toHaveBeenCalled(); + expect(onComplete).toHaveBeenCalledTimes(1); + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("drops the restore when Run clears the canvas before the fetches return", () => { + // Run calls resetExecutionState() before the backend answers, so the execution state is + // briefly Uninitialized and isExecuting() reads false. resetStatus() writes the cleared + // statistics first, and that write is what the restore has to yield to. + executionsService.retrieveWorkflowRuntimeStatistics.mockImplementation(() => { + statisticsUpdates.next({}); + return of([makeStatsRow({})]); + }); + + service.restoreLatestRunStatistics().subscribe(); + + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); + + it("still restores when nothing else writes statistics first", () => { + service.restoreLatestRunStatistics().subscribe(); + + expect(statusService.setExternalStatus).toHaveBeenCalledTimes(1); + }); + + it("surfaces an ingestion failure instead of swallowing it as a skipped restore", () => { + statusService.setExternalStatus.mockImplementation(() => { + throw new Error("ingest failed"); + }); + const onError = vi.fn(); + + service.restoreLatestRunStatistics().subscribe({ error: onError }); + + expect(onError).toHaveBeenCalledWith(expect.objectContaining({ message: "ingest failed" })); + }); + + it("does not ingest an empty statistics payload", () => { + executionsService.retrieveWorkflowRuntimeStatistics.mockReturnValue(of([])); + + service.restoreLatestRunStatistics().subscribe(); + + expect(statusService.setExternalStatus).not.toHaveBeenCalled(); + }); +}); diff --git a/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.ts b/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.ts new file mode 100644 index 0000000000..b541e131e5 --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/heatmap-stats-restore.service.ts @@ -0,0 +1,124 @@ +/** + * 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 { Injectable } from "@angular/core"; +import { EMPTY, Observable, defer } from "rxjs"; +import { catchError, map, switchMap, takeUntil, tap } from "rxjs/operators"; +import { WorkflowExecutionsService } from "../../../dashboard/service/user/workflow-executions/workflow-executions.service"; +import { EXECUTION_STATUS_CODE, WorkflowExecutionsEntry } from "../../../dashboard/type/workflow-executions-entry"; +import { WorkflowActionService } from "../workflow-graph/model/workflow-action.service"; +import { WorkflowStatusService } from "../workflow-status/workflow-status.service"; +import { ExecuteWorkflowService } from "../execute-workflow/execute-workflow.service"; +import { ExecutionState, isNotInExecution } from "../../types/execute-workflow.interface"; +import { loadPersistedHeatmapView } from "./heatmap-overlay-persistence"; +import { toOperatorRuntimeStatusMap } from "./runtime-statistics-mapper"; + +/** + * Restores the last execution's per-operator statistics after a page refresh, + * so the performance heat-map can render a finished run without re-executing. + * + * All gating lives here rather than in the caller: the persisted-overlay check + * is a synchronous localStorage read, so a user who never enabled the overlay + * incurs no fetch at all. + */ +@Injectable({ + providedIn: "root", +}) +export class HeatmapStatsRestoreService { + constructor( + private workflowExecutionsService: WorkflowExecutionsService, + private workflowActionService: WorkflowActionService, + private workflowStatusService: WorkflowStatusService, + private executeWorkflowService: ExecuteWorkflowService + ) {} + + /** + * Fetches the latest run's statistics and feeds them into + * WorkflowStatusService. Cold: nothing happens until subscribed. Skips + * silently (including on HTTP errors — restoring is best-effort) when: + * - the overlay is not persisted on, + * - the workflow has never been saved (no wid), + * - an execution is in progress, on entry or by the time the fetches return + * (the live stream wins), + * - the workflow has no executions or the run left no statistics, + * - another producer writes statistics first (a new run clears the canvas). + */ + public restoreLatestRunStatistics(): Observable<void> { + return defer(() => { + if (loadPersistedHeatmapView() === null) { + return EMPTY; + } + const wid = this.workflowActionService.getWorkflowMetadata()?.wid; + if (wid === undefined) { + return EMPTY; + } + if (this.isExecuting()) { + return EMPTY; + } + + return this.workflowExecutionsService.retrieveWorkflowExecutions(wid).pipe( + switchMap(executions => { + const run = this.pickLatestRun(executions); + if (run === undefined) { + return EMPTY; + } + return this.workflowExecutionsService.retrieveWorkflowRuntimeStatistics(wid, run.eId, run.cuId); + }), + // Only the fetches are best-effort. Placed above the map so a mapping or ingestion + // failure still surfaces instead of looking like a run with nothing to restore. + catchError(() => EMPTY), + tap(rows => { + const runtimeStatus = toOperatorRuntimeStatusMap(rows); + // Re-checked, not redundant: the entry guard ran two round trips ago, and the + // websocket can connect and start streaming a live run inside that window. + if (this.isExecuting() || Object.keys(runtimeStatus).length === 0) { + return; + } + this.workflowStatusService.setExternalStatus(runtimeStatus); + }), + map(() => undefined), + // Any other producer writing statistics means the canvas is no longer ours to restore: + // pressing Run resets the execution state to Uninitialized, which isExecuting() cannot + // see until the backend answers, but resetStatus() writes here first. Unsubscribing + // tears the pending fetch down, so the tap above never runs. A plain Subject, so + // subscribing does not itself emit. The restore's own write does fire it, from inside + // the tap, which closes the stream only after that write has reached its subscribers. + takeUntil(this.workflowStatusService.getStatisticsUpdateStream()) + ); + }); + } + + /** + * The run to restore: the most recent completed execution, or — when no run + * ever completed — the most recent one overall, so a partially executed + * workflow still shows the statistics it produced. + */ + private pickLatestRun(executions: ReadonlyArray<WorkflowExecutionsEntry>): WorkflowExecutionsEntry | undefined { + const latestOf = (entries: ReadonlyArray<WorkflowExecutionsEntry>) => + entries.length === 0 + ? undefined + : entries.reduce((latest, entry) => (entry.startingTime > latest.startingTime ? entry : latest)); + const completed = executions.filter(e => EXECUTION_STATUS_CODE[e.status] === ExecutionState.Completed); + return latestOf(completed) ?? latestOf(executions); + } + + private isExecuting(): boolean { + return !isNotInExecution(this.executeWorkflowService.getExecutionState().state); + } +} diff --git a/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.spec.ts b/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.spec.ts new file mode 100644 index 0000000000..64d1c0bf83 --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.spec.ts @@ -0,0 +1,130 @@ +/** + * 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 { operatorStateFromStatusCode, toOperatorRuntimeStatusMap } from "./runtime-statistics-mapper"; +import { OperatorState } from "../../types/execute-workflow.interface"; +import { WorkflowRuntimeStatistics } from "../../../dashboard/type/workflow-runtime-statistics"; + +function makeRow(overrides: Partial<WorkflowRuntimeStatistics>): WorkflowRuntimeStatistics { + return { + operatorId: "op1", + timestamp: 1_000, + inputTupleCount: 0, + inputTupleSize: 0, + outputTupleCount: 0, + outputTupleSize: 0, + totalDataProcessingTime: 0, + totalControlProcessingTime: 0, + totalIdleTime: 0, + numberOfWorkers: 1, + status: 3, + ...overrides, + }; +} + +describe("operatorStateFromStatusCode", () => { + it("maps the persisted codes the engine writes (Utils.maptoStatusCode)", () => { + expect(operatorStateFromStatusCode(0)).toBe(OperatorState.Uninitialized); + expect(operatorStateFromStatusCode(1)).toBe(OperatorState.Running); + expect(operatorStateFromStatusCode(2)).toBe(OperatorState.Paused); + expect(operatorStateFromStatusCode(3)).toBe(OperatorState.Completed); + }); + + it("falls back to Uninitialized for codes OperatorState cannot express", () => { + // 4 = Failed, 5 = Killed, -1 = other: no OperatorState member exists for + // them, so they render as the neutral default rather than a wrong state. + expect(operatorStateFromStatusCode(4)).toBe(OperatorState.Uninitialized); + expect(operatorStateFromStatusCode(5)).toBe(OperatorState.Uninitialized); + expect(operatorStateFromStatusCode(-1)).toBe(OperatorState.Uninitialized); + expect(operatorStateFromStatusCode(999)).toBe(OperatorState.Uninitialized); + }); +}); + +describe("toOperatorRuntimeStatusMap", () => { + it("maps every field of a row onto the wire shape", () => { + const result = toOperatorRuntimeStatusMap([ + makeRow({ + operatorId: "op1", + inputTupleCount: 1_000, + inputTupleSize: 8_000, + outputTupleCount: 250, + outputTupleSize: 2_000, + totalDataProcessingTime: 5_000_000, + totalControlProcessingTime: 1_000_000, + totalIdleTime: 700_000, + numberOfWorkers: 2, + status: 1, + }), + ]); + + expect(result).toEqual({ + op1: { + operatorState: OperatorState.Running, + aggregatedInputRowCount: 1_000, + aggregatedInputSize: 8_000, + aggregatedOutputRowCount: 250, + aggregatedOutputSize: 2_000, + numWorkers: 2, + aggregatedDataProcessingTime: 5_000_000, + aggregatedControlProcessingTime: 1_000_000, + aggregatedIdleTime: 700_000, + }, + }); + }); + + it("omits the port maps rather than emitting empty ones, so port labels survive a restore", () => { + // An empty map reads as "every port measured zero" and makes JointUIService write 0 over + // the port display names; absent means "no per-port information" and leaves them alone. + const restored = toOperatorRuntimeStatusMap([makeRow({ operatorId: "op1" })]).op1; + expect(restored).not.toHaveProperty("inputPortMetrics"); + expect(restored).not.toHaveProperty("outputPortMetrics"); + }); + + it("keeps only the latest row per operator, regardless of input order", () => { + const result = toOperatorRuntimeStatusMap([ + makeRow({ operatorId: "op1", timestamp: 3_000, outputTupleCount: 30, status: 3 }), + makeRow({ operatorId: "op1", timestamp: 1_000, outputTupleCount: 10, status: 1 }), + makeRow({ operatorId: "op2", timestamp: 2_000, outputTupleCount: 99 }), + makeRow({ operatorId: "op1", timestamp: 2_000, outputTupleCount: 20, status: 1 }), + ]); + + expect(Object.keys(result).sort()).toEqual(["op1", "op2"]); + expect(result["op1"].aggregatedOutputRowCount).toBe(30); + expect(result["op1"].operatorState).toBe(OperatorState.Completed); + expect(result["op2"].aggregatedOutputRowCount).toBe(99); + }); + + it("lets the later row win on equal timestamps (snapshots arrive in write order)", () => { + const result = toOperatorRuntimeStatusMap([ + makeRow({ operatorId: "op1", timestamp: 1_000, outputTupleCount: 1 }), + makeRow({ operatorId: "op1", timestamp: 1_000, outputTupleCount: 2 }), + ]); + expect(result["op1"].aggregatedOutputRowCount).toBe(2); + }); + + it("keys the result by operator id, including unicode ids", () => { + const id = "算子-✓-1"; + const result = toOperatorRuntimeStatusMap([makeRow({ operatorId: id })]); + expect(Object.keys(result)).toEqual([id]); + }); + + it("returns an empty map for empty input", () => { + expect(toOperatorRuntimeStatusMap([])).toEqual({}); + }); +}); diff --git a/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.ts b/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.ts new file mode 100644 index 0000000000..980c6bd6da --- /dev/null +++ b/frontend/src/app/workspace/service/heatmap/runtime-statistics-mapper.ts @@ -0,0 +1,78 @@ +/** + * 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 { OperatorRuntimeStatus, OperatorState } from "../../types/execute-workflow.interface"; +import { WorkflowRuntimeStatistics } from "../../../dashboard/type/workflow-runtime-statistics"; + +/** + * Maps a persisted per-operator status code back to an OperatorState. + * + * The engine persists the codes written by Utils.maptoStatusCode: 0 = + * Uninitialized/Ready, 1 = Running, 2 = Paused, 3 = Completed. Codes 4 + * (Failed), 5 (Killed) and -1 (other) have no OperatorState member, so they + * fall back to the neutral Uninitialized rather than a wrong state. + */ +export function operatorStateFromStatusCode(code: number): OperatorState { + switch (code) { + case 0: + return OperatorState.Uninitialized; + case 1: + return OperatorState.Running; + case 2: + return OperatorState.Paused; + case 3: + return OperatorState.Completed; + default: + return OperatorState.Uninitialized; + } +} + +/** + * Reduces a persisted runtime-statistics time series to the latest snapshot + * per operator (by timestamp; the later row wins a tie, matching write + * order), mapped to the wire shape WorkflowStatusService ingests. + * + * The engine sums the ports away before persisting, so the port maps are left + * absent rather than empty — an empty map would zero the port labels. + */ +export function toOperatorRuntimeStatusMap(rows: WorkflowRuntimeStatistics[]): Record<string, OperatorRuntimeStatus> { + const latestByOperator: Record<string, WorkflowRuntimeStatistics> = {}; + for (const row of rows) { + const seen = latestByOperator[row.operatorId]; + if (seen === undefined || row.timestamp >= seen.timestamp) { + latestByOperator[row.operatorId] = row; + } + } + + const result: Record<string, OperatorRuntimeStatus> = {}; + for (const [operatorId, row] of Object.entries(latestByOperator)) { + result[operatorId] = { + operatorState: operatorStateFromStatusCode(row.status), + aggregatedInputRowCount: row.inputTupleCount, + aggregatedInputSize: row.inputTupleSize, + aggregatedOutputRowCount: row.outputTupleCount, + aggregatedOutputSize: row.outputTupleSize, + numWorkers: row.numberOfWorkers, + aggregatedDataProcessingTime: row.totalDataProcessingTime, + aggregatedControlProcessingTime: row.totalControlProcessingTime, + aggregatedIdleTime: row.totalIdleTime, + }; + } + return result; +} diff --git a/frontend/src/app/workspace/service/joint-ui/joint-ui.service.spec.ts b/frontend/src/app/workspace/service/joint-ui/joint-ui.service.spec.ts index 6a7a4806b2..eeeab17390 100644 --- a/frontend/src/app/workspace/service/joint-ui/joint-ui.service.spec.ts +++ b/frontend/src/app/workspace/service/joint-ui/joint-ui.service.spec.ts @@ -814,5 +814,40 @@ describe("JointUIService", () => { const valuesWritten = attrSpy.mock.calls.map(c => c[1]); expect(valuesWritten).toContain("#workers: 1"); }); + + it("zeroes the port labels for an empty metrics map", () => { + const { paper, portPropSpy } = makeStatsPaper(() => [ + { id: "in-0", group: "in", attrs: { ".port-label": { text: "data" } } }, + { id: "out-1", group: "out", attrs: { ".port-label": { text: "result" } } }, + ]); + const service = new JointUIService(emptyMetadataStub as never); + service.changeOperatorStatistics(paper, "op-1", { + aggregatedInputRowCount: 0, + aggregatedOutputRowCount: 0, + inputPortMetrics: {}, + outputPortMetrics: {}, + }); + expect(portPropSpy).toHaveBeenCalledWith("in-0", "attrs/.port-label/text", (0).toLocaleString()); + expect(portPropSpy).toHaveBeenCalledWith("out-1", "attrs/.port-label/text", (0).toLocaleString()); + }); + + it("leaves the port labels untouched when the metrics maps are absent", () => { + // Statistics restored from a finished run carry no per-port information; writing 0 would + // overwrite the port display names and stick, because the editor reapplies this snapshot + // on operator-add. + const { paper, portPropSpy, attrSpy } = makeStatsPaper(() => [ + { id: "in-0", group: "in", attrs: { ".port-label": { text: "data" } } }, + { id: "out-1", group: "out", attrs: { ".port-label": { text: "result" } } }, + ]); + const service = new JointUIService(emptyMetadataStub as never); + service.changeOperatorStatistics(paper, "op-1", { + aggregatedInputRowCount: 1_000, + aggregatedOutputRowCount: 250, + numWorkers: 2, + }); + expect(portPropSpy).not.toHaveBeenCalled(); + // The rest of the snapshot still renders. + expect(attrSpy.mock.calls.map(c => c[1])).toContain("#workers: 2"); + }); }); }); diff --git a/frontend/src/app/workspace/service/joint-ui/joint-ui.service.ts b/frontend/src/app/workspace/service/joint-ui/joint-ui.service.ts index bb3ba9e57f..4ad263c2a6 100644 --- a/frontend/src/app/workspace/service/joint-ui/joint-ui.service.ts +++ b/frontend/src/app/workspace/service/joint-ui/joint-ui.service.ts @@ -392,25 +392,31 @@ export class JointUIService { const workerCount = statistics.numWorkers ?? 1; element.attr(`.${operatorWorkerCountClass}/text`, "#workers: " + String(workerCount)); - inPorts.forEach(portDef => { - const portId = portDef.id; - if (portId != null) { - const parts = portId.split("-"); - const numericSuffix = parts.length > 1 ? parts[1] : portId; - const count: number = inputMetrics[numericSuffix] ?? 0; - element.portProp(portId, "attrs/.port-label/text", count.toLocaleString()); - } - }); + // Absent map: no per-port information in this snapshot, so leave the labels alone. + // Empty map: every port measured zero, so write the zeros. + if (inputMetrics !== undefined) { + inPorts.forEach(portDef => { + const portId = portDef.id; + if (portId != null) { + const parts = portId.split("-"); + const numericSuffix = parts.length > 1 ? parts[1] : portId; + const count: number = inputMetrics[numericSuffix] ?? 0; + element.portProp(portId, "attrs/.port-label/text", count.toLocaleString()); + } + }); + } - outPorts.forEach(portDef => { - const portId = portDef.id; - if (portId != null) { - const parts = portId.split("-"); - const numericSuffix = parts.length > 1 ? parts[1] : portId; - const count: number = outputMetrics[numericSuffix] ?? 0; - element.portProp(portId, "attrs/.port-label/text", count.toLocaleString()); - } - }); + if (outputMetrics !== undefined) { + outPorts.forEach(portDef => { + const portId = portDef.id; + if (portId != null) { + const parts = portId.split("-"); + const numericSuffix = parts.length > 1 ? parts[1] : portId; + const count: number = outputMetrics[numericSuffix] ?? 0; + element.portProp(portId, "attrs/.port-label/text", count.toLocaleString()); + } + }); + } } public foldOperatorDetails(jointPaper: joint.dia.Paper, operatorID: string): void { jointPaper.getModelById(operatorID).attr({ diff --git a/frontend/src/app/workspace/service/workflow-status/workflow-status.service.spec.ts b/frontend/src/app/workspace/service/workflow-status/workflow-status.service.spec.ts index 7aa5b6aa53..66c4a4241f 100644 --- a/frontend/src/app/workspace/service/workflow-status/workflow-status.service.spec.ts +++ b/frontend/src/app/workspace/service/workflow-status/workflow-status.service.spec.ts @@ -63,6 +63,19 @@ describe("WorkflowStatusService", () => { service = TestBed.inject(WorkflowStatusService); }); + it("emits nothing on subscribe, before any producer has written", () => { + // Load-bearing for HeatmapStatsRestoreService, which uses this stream as a takeUntil + // notifier: were it a BehaviorSubject, the restore would be cancelled on every page load. + const stateEmissions: Record<string, OperatorState>[] = []; + const statisticsEmissions: Record<string, OperatorStatistics>[] = []; + + service.getStateUpdateStream().subscribe(s => stateEmissions.push(s)); + service.getStatisticsUpdateStream().subscribe(s => statisticsEmissions.push(s)); + + expect(stateEmissions).toHaveLength(0); + expect(statisticsEmissions).toHaveLength(0); + }); + it("splits an OperatorStatisticsUpdateEvent into the state and statistics streams", () => { const stateEmissions: Record<string, OperatorState>[] = []; const statisticsEmissions: Record<string, OperatorStatistics>[] = []; @@ -220,4 +233,27 @@ describe("WorkflowStatusService", () => { expect(service.getCurrentStatistics()).toEqual({}); expect(service.getCurrentPerformanceMetrics()).toEqual({}); }); + + it("setExternalStatus ingests through the same split path as a websocket update", () => { + const order: string[] = []; + service.getStateUpdateStream().subscribe(() => order.push("state")); + service.getStatisticsUpdateStream().subscribe(() => order.push("statistics")); + + service.setExternalStatus({ op1: sampleRuntimeStatus }); + + // Same contract as the wire path: state first, statistics second, no + // operatorState leaking into the statistics concept, metrics derived. + expect(order).toEqual(["state", "statistics"]); + expect(service.getCurrentState()).toEqual({ op1: OperatorState.Running }); + expect(service.getCurrentStatistics()).toEqual({ op1: sampleStatistics }); + expect(service.getCurrentStatistics()["op1"]).not.toHaveProperty("operatorState"); + expect(service.getCurrentPerformanceMetrics()["op1"].dataProcessingTimeNs).toBe(5_000_000); + }); + + it("a later live websocket update overrides externally restored status", () => { + service.setExternalStatus({ op1: sampleRuntimeStatus }); + websocketEventSubject.next(statsEvent({ op1: { ...sampleRuntimeStatus, operatorState: OperatorState.Paused } })); + + expect(service.getCurrentState()).toEqual({ op1: OperatorState.Paused }); + }); }); diff --git a/frontend/src/app/workspace/service/workflow-status/workflow-status.service.ts b/frontend/src/app/workspace/service/workflow-status/workflow-status.service.ts index 729adb6507..0e20939d1b 100644 --- a/frontend/src/app/workspace/service/workflow-status/workflow-status.service.ts +++ b/frontend/src/app/workspace/service/workflow-status/workflow-status.service.ts @@ -19,7 +19,7 @@ import { Injectable } from "@angular/core"; import { BehaviorSubject, Observable, Subject } from "rxjs"; -import { OperatorState, OperatorStatistics } from "../../types/execute-workflow.interface"; +import { OperatorRuntimeStatus, OperatorState, OperatorStatistics } from "../../types/execute-workflow.interface"; import { WorkflowWebsocketService } from "../workflow-websocket/workflow-websocket.service"; import { OperatorPerformanceMetrics, extractPerformanceMetrics } from "./performance-metrics"; @@ -64,18 +64,33 @@ export class WorkflowStatusService { if (event.type !== "OperatorStatisticsUpdateEvent") { return; } - const state: Record<string, OperatorState> = {}; - const statistics: Record<string, OperatorStatistics> = {}; - for (const [operatorId, update] of Object.entries(event.operatorStatistics)) { - const { operatorState, ...statisticsOnly } = update; - state[operatorId] = operatorState; - statistics[operatorId] = statisticsOnly; - } - this.stateSubject.next(state); - this.statisticsSubject.next(statistics); + this.ingestRuntimeStatus(event.operatorStatistics); }); } + /** Splits a bundled runtime-status map into the two sub-concept emissions. */ + private ingestRuntimeStatus(runtimeStatus: Record<string, OperatorRuntimeStatus>): void { + const state: Record<string, OperatorState> = {}; + const statistics: Record<string, OperatorStatistics> = {}; + for (const [operatorId, update] of Object.entries(runtimeStatus)) { + const { operatorState, ...statisticsOnly } = update; + state[operatorId] = operatorState; + statistics[operatorId] = statisticsOnly; + } + this.stateSubject.next(state); + this.statisticsSubject.next(statistics); + } + + /** + * Ingests externally sourced per-operator runtime status (e.g. a finished + * run's persisted statistics restored after a page refresh), through the + * same split-and-emit path as live websocket updates — the emission order + * and derived performance metrics behave identically. + */ + public setExternalStatus(runtimeStatus: Record<string, OperatorRuntimeStatus>): void { + this.ingestRuntimeStatus(runtimeStatus); + } + /** Stream of per-operator execution states, keyed by operator id. */ public getStateUpdateStream(): Observable<Record<string, OperatorState>> { return this.stateSubject.asObservable(); diff --git a/frontend/src/app/workspace/types/execute-workflow.interface.ts b/frontend/src/app/workspace/types/execute-workflow.interface.ts index baee2d6594..79f8aa8d2b 100644 --- a/frontend/src/app/workspace/types/execute-workflow.interface.ts +++ b/frontend/src/app/workspace/types/execute-workflow.interface.ts @@ -82,10 +82,11 @@ export interface OperatorStatistics extends Readonly<{ aggregatedInputRowCount: number; aggregatedInputSize?: number; - inputPortMetrics: Record<string, number>; + /** Absent when the snapshot has no per-port information at all; `{}` means every port was zero. */ + inputPortMetrics?: Record<string, number>; aggregatedOutputRowCount: number; aggregatedOutputSize?: number; - outputPortMetrics: Record<string, number>; + outputPortMetrics?: Record<string, number>; numWorkers?: number; aggregatedDataProcessingTime?: number; aggregatedControlProcessingTime?: number;
