This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/texera.git
The following commit(s) were added to refs/heads/main by this push:
new f9b899a528 refactor(frontend): split operator state from statistics
(#8301)
f9b899a528 is described below
commit f9b899a52890347d8d44c72c325863ff96fcb714
Author: Prateek Ganigi <[email protected]>
AuthorDate: Wed Sep 16 04:01:43 2026 +0000
refactor(frontend): split operator state from statistics (#8301)
### What changes were proposed in this PR?
`WorkflowStatusService` currently bundles two different concepts in one
object: `OperatorStatistics` carries both the operator's execution
**state** (Running, Completed, …) and its **statistics** (row counts,
sizes, timing). This PR splits them into separate sub-concepts, so the
service now exposes three cleanly separated things: state, statistics,
and performance metrics (the third was already separate, from #5834).
- `WorkflowStatusService` now has a stream + snapshot pair per concept:
`getStateUpdateStream()` / `getCurrentState()` for state, and
`getStatisticsUpdateStream()` / `getCurrentStatistics()` for statistics
(metrics only). The performance-metrics API is unchanged.
- `OperatorStatistics` no longer contains `operatorState`. The combined
shape the engine still sends over the websocket is typed as
`OperatorRuntimeStatus`, and the service splits each update into the two
maps. No backend or wire-format changes.
- All consumers are migrated: components that only cared about state
(result panel, code debugger, UDF debug service, property editor) now
read the state stream; the workflow editor renders state (operator
color) and statistics (port counts, worker count) from their own
streams. `JointUIService.changeOperatorStatistics` renders statistics
only — state rendering stays in `changeOperatorState` (its two
long-unused `isSource`/`isSink` params are dropped along the way).
- A small `WorkflowGraph.getAllOperatorIDs()` accessor keeps the
per-update rendering path from materializing full operator predicates
when only IDs are needed.
The change is behavior-preserving. One deliberate exception: the old
code applied the "Recovering" display state by mutating the shared
emitted map, which leaked masked states to other subscribers depending
on subscription order. That accident is removed, and the override is now
applied explicitly where state is rendered.
Rebased on top of the merged heat-map overlay (#6213), which consumes
only the unchanged performance-metrics stream; its editor wiring is
untouched by this refactor.
### Any related issues, documentation, discussions?
Closes #5919. Part of umbrella #5772. Follow-up from the review
discussion in #5834; follows RFC discussion #5216.
### How was this PR tested?
The `WorkflowStatusService` spec now asserts state and statistics are
exposed and update independently, and that statistics never leak
`operatorState`. Consumer specs were updated to the new API, plus new
tests for the state-rendering rules in the workflow editor
(Uninitialized fallback, Recovering override, state label restored after
navigation).
Full frontend suite passes (5,350 tests, 209 files); `tsc --noEmit`,
`eslint ./src`, and Prettier are all clean.
### Was this PR authored or co-authored using generative AI tooling?
This PR was co-authored using Claude in compliance with ASF.
---------
Co-authored-by: Claude Fable 5 <[email protected]>
---
.../code-debugger.component.spec.ts | 74 ++-------
.../code-editor-dialog/code-debugger.component.ts | 10 +-
.../operator-property-edit-frame.component.html | 2 +-
.../operator-property-edit-frame.component.spec.ts | 22 ++-
.../operator-property-edit-frame.component.ts | 10 +-
.../result-table-frame.component.spec.ts | 20 +--
.../result-table-frame.component.ts | 6 +-
.../workflow-editor.component.spec.ts | 180 ++++++++++++---------
.../workflow-editor/workflow-editor.component.ts | 156 +++++++++---------
.../service/joint-ui/joint-ui.service.spec.ts | 87 +++++-----
.../workspace/service/joint-ui/joint-ui.service.ts | 13 +-
.../operator-debug/udf-debug.service.spec.ts | 33 ++--
.../service/operator-debug/udf-debug.service.ts | 4 +-
.../service/workflow-graph/model/workflow-graph.ts | 9 ++
.../workflow-status/performance-metrics.spec.ts | 4 +-
.../service/workflow-status/performance-metrics.ts | 9 +-
.../workflow-status.service.spec.ts | 116 ++++++++++---
.../workflow-status/workflow-status.service.ts | 101 ++++++++----
.../workspace/types/execute-workflow.interface.ts | 15 +-
19 files changed, 471 insertions(+), 400 deletions(-)
diff --git
a/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.spec.ts
b/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.spec.ts
index dbbfe40e10..e4cc086ffc 100644
---
a/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.spec.ts
+++
b/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.spec.ts
@@ -25,7 +25,7 @@ import { UdfDebugService } from
"../../service/operator-debug/udf-debug.service"
import { Subject } from "rxjs";
import * as Y from "yjs";
import { BreakpointInfo } from "../../types/workflow-common.interface";
-import { OperatorState, OperatorStatistics } from
"../../types/execute-workflow.interface";
+import { OperatorState } from "../../types/execute-workflow.interface";
import { commonTestProviders } from "../../../common/testing/test-utils";
import type { Mocked } from "vitest";
import type { MonacoBreakpoint } from "monaco-breakpoints";
@@ -38,14 +38,14 @@ describe("CodeDebuggerComponent", () => {
let mockWorkflowStatusService: Mocked<WorkflowStatusService>;
let mockUdfDebugService: Mocked<UdfDebugService>;
- let statusUpdateStream: Subject<Record<string, OperatorStatistics>>;
+ let stateUpdateStream: Subject<Record<string, OperatorState>>;
let debugState: Y.Map<BreakpointInfo>;
const operatorId = "test-operator-id";
beforeEach(async () => {
// Initialize streams and spy objects
- statusUpdateStream = new Subject<Record<string, OperatorStatistics>>();
+ stateUpdateStream = new Subject<Record<string, OperatorState>>();
// Y.Map observers only fire when the map is attached to a Y.Doc (the doc
// owns the transaction lifecycle that drives observation). A standalone
// `new Y.Map()` accepts `.set()` but never notifies observers — production
@@ -53,8 +53,8 @@ describe("CodeDebuggerComponent", () => {
// doc, but the spec used to construct a detached map.
debugState = new Y.Doc().getMap<BreakpointInfo>("debug");
- mockWorkflowStatusService = { getStatusUpdateStream: vi.fn() } as unknown
as Mocked<WorkflowStatusService>;
-
mockWorkflowStatusService.getStatusUpdateStream.mockReturnValue(statusUpdateStream.asObservable());
+ mockWorkflowStatusService = { getStateUpdateStream: vi.fn() } as unknown
as Mocked<WorkflowStatusService>;
+
mockWorkflowStatusService.getStateUpdateStream.mockReturnValue(stateUpdateStream.asObservable());
mockUdfDebugService = {
getDebugState: vi.fn(),
@@ -85,7 +85,7 @@ describe("CodeDebuggerComponent", () => {
afterEach(() => {
// Clean up streams to prevent memory leaks
- statusUpdateStream.complete();
+ stateUpdateStream.complete();
component.monacoEditor?.dispose();
});
@@ -103,15 +103,7 @@ describe("CodeDebuggerComponent", () => {
const rerenderSpy = vi.spyOn(component,
"rerenderExistingBreakpoints").mockImplementation(() => {});
// Emit a Running state event
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Running,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Running });
tick();
fixture.detectChanges(); // Trigger change detection
@@ -120,15 +112,7 @@ describe("CodeDebuggerComponent", () => {
expect(rerenderSpy).toHaveBeenCalled();
// Emit the same state again (should not trigger setup again)
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Running,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Running });
tick();
fixture.detectChanges(); // Trigger change detection
@@ -137,15 +121,7 @@ describe("CodeDebuggerComponent", () => {
expect(rerenderSpy).toHaveBeenCalledTimes(1); // No additional call
// Emit the paused state (should not trigger setup)
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Paused,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Paused });
tick();
fixture.detectChanges(); // Trigger change detection
@@ -154,15 +130,7 @@ describe("CodeDebuggerComponent", () => {
expect(rerenderSpy).toHaveBeenCalledTimes(1); // No additional call
// Emit the running state once more (should not trigger setup)
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Paused,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Paused });
tick();
fixture.detectChanges(); // Trigger change detection
@@ -175,30 +143,14 @@ describe("CodeDebuggerComponent", () => {
const removeSpy = vi.spyOn(component, "removeMonacoBreakpointMethods");
// Emit an Uninitialized state event
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Uninitialized,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Uninitialized });
fixture.detectChanges(); // Trigger change detection
expect(removeSpy).toHaveBeenCalled();
// Emit the same state again (should not trigger removal again)
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Uninitialized,
- aggregatedOutputRowCount: 0,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Uninitialized });
expect(removeSpy).toHaveBeenCalledTimes(1); // No additional call
});
@@ -499,7 +451,7 @@ describe("CodeDebuggerComponent breakpoint gutter", () => {
providers: [
{
provide: WorkflowStatusService,
- useValue: { getStatusUpdateStream: vi.fn(() => new
Subject().asObservable()) },
+ useValue: { getStateUpdateStream: vi.fn(() => new
Subject().asObservable()) },
},
{ provide: UdfDebugService, useValue: debugService },
...commonTestProviders,
diff --git
a/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.ts
b/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.ts
index d46b19f734..9c5ebacfbf 100644
---
a/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.ts
+++
b/frontend/src/app/workspace/component/code-editor-dialog/code-debugger.component.ts
@@ -59,7 +59,7 @@ export class CodeDebuggerComponent implements AfterViewInit,
SafeStyle {
) {}
ngAfterViewInit() {
- this.registerStatusChangeHandler();
+ this.registerStateChangeHandler();
this.registerBreakpointRenderingHandler();
}
@@ -221,14 +221,14 @@ export class CodeDebuggerComponent implements
AfterViewInit, SafeStyle {
});
}
- private registerStatusChangeHandler() {
+ private registerStateChangeHandler() {
this.workflowStatusService
- .getStatusUpdateStream()
+ .getStateUpdateStream()
.pipe(
map(
event =>
- event[this.currentOperatorId]?.operatorState ===
OperatorState.Running ||
- event[this.currentOperatorId]?.operatorState ===
OperatorState.Paused
+ event[this.currentOperatorId] === OperatorState.Running ||
+ event[this.currentOperatorId] === OperatorState.Paused
),
distinctUntilChanged(),
untilDestroyed(this)
diff --git
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.html
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.html
index 4e98e35d97..5554faeea2 100644
---
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.html
+++
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.html
@@ -213,7 +213,7 @@
nzTooltipPlacement="bottom"
[disabled]="
currentOperatorSchema?.additionalMetadata?.supportReconfiguration !== true
- || currentOperatorStatus?.operatorState === OperatorState.Completed
+ || currentOperatorState === OperatorState.Completed
">
Unlock for Logic Change
<i
diff --git
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
index d8e9cedeb7..fa551d1402 100644
---
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
+++
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.spec.ts
@@ -60,6 +60,7 @@ import { WorkflowGraph } from
"../../../service/workflow-graph/model/workflow-gr
import { UiUdfParametersSyncService } from
"../../../service/code-editor/ui-udf-parameters-sync.service";
import { WorkflowPveService } from
"../../../service/virtual-environment/virtual-environment.service";
import { WorkflowWebsocketService } from
"../../../service/workflow-websocket/workflow-websocket.service";
+import { OperatorState } from "../../../types/execute-workflow.interface";
import { TexeraWebsocketEvent } from
"../../../types/workflow-websocket.interface";
import { of, Subject, throwError } from "rxjs";
import { WorkflowVersionService } from
"../../../../dashboard/service/user/workflow-version/workflow-version.service";
@@ -2681,23 +2682,32 @@ describe("OperatorPropertyEditFrameComponent", () => {
expect(rerenderSpy).not.toHaveBeenCalled();
});
- it("the status-update subscription records the update for the selected
operator", () => {
+ // Wire-shaped payload: the service splits it into the state and
statistics concepts.
+ const runningStatus = {
+ operatorState: OperatorState.Running,
+ aggregatedInputRowCount: 0,
+ inputPortMetrics: {},
+ aggregatedOutputRowCount: 0,
+ outputPortMetrics: {},
+ };
+
+ it("the state-update subscription records the state for the selected
operator", () => {
workflowActionService.addOperator(mockScanPredicate, mockPoint);
component.currentOperatorId = mockScanPredicate.operatorID;
fixture.detectChanges(); // ngOnInit registers the subscription
- emitOperatorStatistics({ [mockScanPredicate.operatorID]: { some:
"status" } });
+ emitOperatorStatistics({ [mockScanPredicate.operatorID]: runningStatus
});
- expect(component.currentOperatorStatus).toEqual({ some: "status" });
+ expect(component.currentOperatorState).toBe(OperatorState.Running);
});
- it("the status-update subscription ignores updates while no operator is
selected", () => {
+ it("the state-update subscription ignores updates while no operator is
selected", () => {
fixture.detectChanges(); // ngOnInit registers the subscription
component.currentOperatorId = undefined;
- emitOperatorStatistics({ "op-1": { some: "status" } });
+ emitOperatorStatistics({ "op-1": runningStatus });
- expect(component.currentOperatorStatus).toBeUndefined();
+ expect(component.currentOperatorState).toBeUndefined();
});
it("the ui-parameter subscription ignores events addressed to a different
operator", () => {
diff --git
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
index b9c4bf5f6d..033b75bc48 100644
---
a/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
+++
b/frontend/src/app/workspace/component/property-editor/operator-property-edit-frame/operator-property-edit-frame.component.ts
@@ -39,7 +39,7 @@ import {
} from "../../../types/custom-json-schema.interface";
import { isDefined } from "../../../../common/util/predicate";
import { customFormlyFieldType, NON_FORM_FIELD_TYPES } from
"../../../util/custom-formly-type";
-import { ExecutionState, OperatorState, OperatorStatistics } from
"src/app/workspace/types/execute-workflow.interface";
+import { ExecutionState, OperatorState } from
"src/app/workspace/types/execute-workflow.interface";
import { DynamicSchemaService } from
"../../../service/dynamic-schema/dynamic-schema.service";
import { WorkflowCompilingService } from
"../../../service/compile-workflow/workflow-compiling.service";
import {
@@ -193,7 +193,7 @@ export class OperatorPropertyEditFrameComponent implements
OnInit, OnChanges, On
currentOperatorSchema?: OperatorSchema;
readonly OperatorState = OperatorState;
- currentOperatorStatus?: OperatorStatistics;
+ currentOperatorState?: OperatorState;
// re-declare enum for angular template to access it
readonly ExecutionState = ExecutionState;
@@ -570,11 +570,11 @@ export class OperatorPropertyEditFrameComponent
implements OnInit, OnChanges, On
this.registerOperatorDisplayNameChangeHandler();
this.workflowStatusSerivce
- .getStatusUpdateStream()
+ .getStateUpdateStream()
.pipe(untilDestroyed(this))
.subscribe(update => {
if (this.currentOperatorId) {
- this.currentOperatorStatus = update[this.currentOperatorId];
+ this.currentOperatorState = update[this.currentOperatorId];
}
});
@@ -644,7 +644,7 @@ export class OperatorPropertyEditFrameComponent implements
OnInit, OnChanges, On
return;
}
this.currentOperatorSchema =
this.dynamicSchemaService.getDynamicSchema(this.currentOperatorId);
- this.currentOperatorStatus =
this.workflowStatusSerivce.getCurrentStatus()[this.currentOperatorId];
+ this.currentOperatorState =
this.workflowStatusSerivce.getCurrentState()[this.currentOperatorId];
if (this.actsAsEditor) {
this.workflowActionService
diff --git
a/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.spec.ts
b/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.spec.ts
index d4c9e85b08..50e7c224a2 100644
---
a/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.spec.ts
+++
b/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.spec.ts
@@ -38,7 +38,7 @@ import {
} from "../../../service/workflow-result/workflow-result.service";
import { WorkflowStatusService } from
"../../../service/workflow-status/workflow-status.service";
import { PanelResizeService } from
"../../../service/workflow-result/panel-resize/panel-resize.service";
-import { OperatorState, OperatorStatistics, WebResultUpdate } from
"../../../types/execute-workflow.interface";
+import { OperatorState, WebResultUpdate } from
"../../../types/execute-workflow.interface";
import { PaginatedResultEvent } from
"../../../types/workflow-websocket.interface";
import { IndexableObject } from "../../../types/result-table.interface";
import { RowModalComponent } from "../result-panel-modal.component";
@@ -84,14 +84,6 @@ describe("ResultTableFrameComponent", () => {
return paginatedResultService;
};
- const makeStatistics = (state: OperatorState): OperatorStatistics => ({
- operatorState: state,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- aggregatedOutputRowCount: 0,
- outputPortMetrics: {},
- });
-
const paginationUpdate = (totalNumTuples: number, dirtyPageIndices:
number[]): WebResultUpdate => ({
mode: { type: "PaginationMode" },
totalNumTuples,
@@ -278,17 +270,17 @@ describe("ResultTableFrameComponent", () => {
describe("workflow status stream", () => {
it("marks the operator finished only while its reported state is
Completed", () => {
- const statusStream = new Subject<Record<string, OperatorStatistics>>();
- vi.spyOn(workflowStatusService,
"getStatusUpdateStream").mockReturnValue(statusStream.asObservable());
+ const stateStream = new Subject<Record<string, OperatorState>>();
+ vi.spyOn(workflowStatusService,
"getStateUpdateStream").mockReturnValue(stateStream.asObservable());
recreateComponent("op1");
- statusStream.next({ op1: makeStatistics(OperatorState.Completed) });
+ stateStream.next({ op1: OperatorState.Completed });
expect(component.isOperatorFinished).toBe(true);
- statusStream.next({ op1: makeStatistics(OperatorState.Running) });
+ stateStream.next({ op1: OperatorState.Running });
expect(component.isOperatorFinished).toBe(false);
- statusStream.next({ otherOp: makeStatistics(OperatorState.Completed) });
+ stateStream.next({ otherOp: OperatorState.Completed });
expect(component.isOperatorFinished).toBe(false);
});
});
diff --git
a/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.ts
b/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.ts
index e46fa49f51..7fcfd126a3 100644
---
a/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.ts
+++
b/frontend/src/app/workspace/component/result-panel/result-table-frame/result-table-frame.component.ts
@@ -151,10 +151,10 @@ export class ResultTableFrameComponent implements OnInit,
OnChanges {
ngOnInit(): void {
this.workflowStatusService
- .getStatusUpdateStream()
+ .getStateUpdateStream()
.pipe(untilDestroyed(this))
- .subscribe(statusMap => {
- if (this.operatorId && statusMap[this.operatorId]?.operatorState ===
OperatorState.Completed) {
+ .subscribe(stateMap => {
+ if (this.operatorId && stateMap[this.operatorId] ===
OperatorState.Completed) {
this.isOperatorFinished = true;
this.changeDetectorRef.detectChanges();
} else {
diff --git
a/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.spec.ts
b/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.spec.ts
index 67111e5476..d397704d10 100644
---
a/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.spec.ts
+++
b/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.spec.ts
@@ -32,6 +32,7 @@ import {
JointUIService,
operatorAgentActionProgressClass,
operatorNameClass,
+ operatorStateClass,
} from "../../service/joint-ui/joint-ui.service";
import { AgentService, OperatorResultSummary } from
"../../service/agent/agent.service";
import { NzModalModule, NzModalService } from "ng-zorro-antd/modal";
@@ -972,23 +973,17 @@ describe("WorkflowEditorComponent", () => {
* default (gray) when the user navigates away from and back to a workflow
* that has already finished executing. Both the operator-add stream and
* the validation stream route their final border decision through
- * applyOperatorBorder, which encodes the priority: invalid > cached
+ * applyOperatorStateAndBorder, which encodes the priority: invalid >
cached
* execution state > default valid. These tests assert the operator's
* actual final rect.body/stroke on the paper, so they pin down the visible
* outcome rather than the internal helper calls.
*/
describe("operator border restoration after navigation", () => {
let workflowStatusService: WorkflowStatusService;
- const cachedStatus = (operatorState: OperatorState) => ({
- [mockScanPredicate.operatorID]: {
- operatorState,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- aggregatedOutputRowCount: 0,
- outputPortMetrics: {},
- },
+ const cachedState = (operatorState: OperatorState) => ({
+ [mockScanPredicate.operatorID]: operatorState,
});
- const cachedCompleted = cachedStatus(OperatorState.Completed);
+ const cachedCompleted = cachedState(OperatorState.Completed);
const getStroke = (operatorID: string): string =>
component.paper.getModelById(operatorID).attr("rect.body/stroke") as
string;
@@ -996,8 +991,8 @@ describe("WorkflowEditorComponent", () => {
workflowStatusService = TestBed.inject(WorkflowStatusService);
});
- it("paints the execution-state stroke (green) for a valid operator with
a cached Completed status", () => {
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue(cachedCompleted);
+ it("paints the execution-state stroke (green) for a valid operator with
a cached Completed state", () => {
+ vi.spyOn(workflowStatusService,
"getCurrentState").mockReturnValue(cachedCompleted);
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: true });
workflowActionService.addOperator(mockScanPredicate, mockPoint);
@@ -1006,10 +1001,10 @@ describe("WorkflowEditorComponent", () => {
expect(getStroke(mockScanPredicate.operatorID)).toBe("green");
});
- it("paints the execution-state stroke (orange) for a valid operator with
a cached Running status", () => {
+ it("paints the execution-state stroke (orange) for a valid operator with
a cached Running state", () => {
// Navigation-return with a mid-run operator: the border must be
restored
// to the running color, not the default (see #3614).
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue(cachedStatus(OperatorState.Running));
+ vi.spyOn(workflowStatusService,
"getCurrentState").mockReturnValue(cachedState(OperatorState.Running));
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: true });
workflowActionService.addOperator(mockScanPredicate, mockPoint);
@@ -1018,8 +1013,8 @@ describe("WorkflowEditorComponent", () => {
expect(getStroke(mockScanPredicate.operatorID)).toBe("orange");
});
- it("falls back to the default valid stroke (#CFCFCF) when no cached
status exists", () => {
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue({});
+ it("falls back to the default valid stroke (#CFCFCF) when no cached
state exists", () => {
+ vi.spyOn(workflowStatusService, "getCurrentState").mockReturnValue({});
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: true });
workflowActionService.addOperator(mockScanPredicate, mockPoint);
@@ -1028,8 +1023,8 @@ describe("WorkflowEditorComponent", () => {
expect(getStroke(mockScanPredicate.operatorID)).toBe("#CFCFCF");
});
- it("paints the invalid stroke (red) for an invalid operator with no
cached status", () => {
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue({});
+ it("paints the invalid stroke (red) for an invalid operator with no
cached state", () => {
+ vi.spyOn(workflowStatusService, "getCurrentState").mockReturnValue({});
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: false, messages: {} });
workflowActionService.addOperator(mockScanPredicate, mockPoint);
@@ -1038,11 +1033,11 @@ describe("WorkflowEditorComponent", () => {
expect(getStroke(mockScanPredicate.operatorID)).toBe("red");
});
- it("prioritizes invalid (red) over cached Completed status", () => {
+ it("prioritizes invalid (red) over cached Completed state", () => {
// Regression case: operator is both invalid AND has a cached Completed
- // status. applyOperatorBorder must pick red regardless of the order in
+ // state. applyOperatorStateAndBorder must pick red regardless of the
order in
// which the operator-add and validation streams fire.
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue(cachedCompleted);
+ vi.spyOn(workflowStatusService,
"getCurrentState").mockReturnValue(cachedCompleted);
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: false, messages: {} });
workflowActionService.addOperator(mockScanPredicate, mockPoint);
@@ -1061,7 +1056,7 @@ describe("WorkflowEditorComponent", () => {
// The helper takes the Validation as a required argument and must use
it
// directly — it has no fallback path that calls validateOperator
itself.
- (component as any).applyOperatorBorder(mockScanPredicate.operatorID, {
isValid: true });
+ (component as
any).applyOperatorStateAndBorder(mockScanPredicate.operatorID, { isValid: true
});
expect(validateSpy).not.toHaveBeenCalled();
});
@@ -1069,30 +1064,71 @@ describe("WorkflowEditorComponent", () => {
it("honors the passed-in Validation result (paints red when it is
invalid)", () => {
// Proves the passed-in value actually drives the border: an invalid
// result must paint red.
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue({});
+ vi.spyOn(workflowStatusService, "getCurrentState").mockReturnValue({});
workflowActionService.addOperator(mockScanPredicate, mockPoint);
fixture.detectChanges();
- (component as any).applyOperatorBorder(mockScanPredicate.operatorID, {
isValid: false, messages: {} });
+ (component as
any).applyOperatorStateAndBorder(mockScanPredicate.operatorID, { isValid:
false, messages: {} });
expect(getStroke(mockScanPredicate.operatorID)).toBe("red");
});
- it("always supplies a Validation to applyOperatorBorder when an operator
is added", () => {
+ it("always supplies a Validation to applyOperatorStateAndBorder when an
operator is added", () => {
// Both subscribers (operator-add and the validation stream) call
- // applyOperatorBorder on add with identical args, so this asserts the
+ // applyOperatorStateAndBorder on add with identical args, so this
asserts the
// required-parameter contract holds through the add flow — every call
// carries a Validation, never undefined — rather than isolating the
// operator-add caller specifically.
- vi.spyOn(workflowStatusService,
"getCurrentStatus").mockReturnValue({});
+ vi.spyOn(workflowStatusService, "getCurrentState").mockReturnValue({});
vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: true });
- const applyBorderSpy = vi.spyOn(component as any,
"applyOperatorBorder");
+ const applyBorderSpy = vi.spyOn(component as any,
"applyOperatorStateAndBorder");
workflowActionService.addOperator(mockScanPredicate, mockPoint);
fixture.detectChanges();
expect(applyBorderSpy).toHaveBeenCalledWith(mockScanPredicate.operatorID, {
isValid: true });
});
+
+ it("restores the execution-state label for an invalid operator with a
cached state", () => {
+ // The red border takes priority for the stroke, but the cached state
+ // must still be rendered (label text), matching how the operator
+ // looked before navigating away: state painted by the state stream,
+ // stroke overridden by validation.
+ vi.spyOn(workflowStatusService,
"getCurrentState").mockReturnValue(cachedCompleted);
+ vi.spyOn(validationWorkflowService,
"validateOperator").mockReturnValue({ isValid: false, messages: {} });
+
+ workflowActionService.addOperator(mockScanPredicate, mockPoint);
+ fixture.detectChanges();
+
+ const stateText = component.paper
+ .getModelById(mockScanPredicate.operatorID)
+ .attr(`.${operatorStateClass}/text`) as string;
+ expect(stateText).toBe(OperatorState.Completed.toString());
+ expect(getStroke(mockScanPredicate.operatorID)).toBe("red");
+ });
+ });
+
+ describe("effectiveOperatorState", () => {
+ const resolve = (reported?: OperatorState): OperatorState => (component
as any).effectiveOperatorState(reported);
+
+ it("falls back to Uninitialized for an operator missing from the state
map", () => {
+ expect(resolve(undefined)).toBe(OperatorState.Uninitialized);
+ });
+
+ it("returns the reported state as-is outside of recovery", () => {
+ expect(resolve(OperatorState.Running)).toBe(OperatorState.Running);
+ });
+
+ it("masks any reported state to Recovering while the execution is
recovering", () => {
+ const executeWorkflowService = TestBed.inject(ExecuteWorkflowService);
+ vi.spyOn(executeWorkflowService, "getExecutionState").mockReturnValue({
+ state: ExecutionState.Recovering,
+ } as ReturnType<ExecuteWorkflowService["getExecutionState"]>);
+
+ expect(resolve(OperatorState.Running)).toBe(OperatorState.Recovering);
+ // the missing-operator fallback is not masked
+ expect(resolve(undefined)).toBe(OperatorState.Uninitialized);
+ });
});
/**
@@ -1842,16 +1878,13 @@ describe("WorkflowEditorComponent editor wiring", () =>
{
(executeWorkflowService as any).regionUpdateStream.next({ regions });
}
- /** An OperatorStatistics payload in the given state. */
- function statisticsIn(state: OperatorState) {
- return {
- operatorState: state,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- aggregatedOutputRowCount: 0,
- outputPortMetrics: {},
- };
- }
+ /** A metrics-only OperatorStatistics payload (operator state travels on its
own stream). */
+ const emptyStatistics = {
+ aggregatedInputRowCount: 0,
+ inputPortMetrics: {},
+ aggregatedOutputRowCount: 0,
+ outputPortMetrics: {},
+ };
/** Clicks the chat button of a cell, the way `.chat-button` does. */
function clickChatButton(cellID: string): void {
@@ -1923,29 +1956,34 @@ describe("WorkflowEditorComponent editor wiring", () =>
{
});
describe("execution status streams", () => {
- it("forwards each operator's statistics, tagging which end of the graph it
sits on", () => {
- workflowActionService.addOperator(mockScanPredicate, mockPoint); // no
input ports -> source
- workflowActionService.addOperator(mockResultPredicate, mockPoint); // no
output ports -> sink
+ it("forwards each operator's statistics from the statistics stream", () =>
{
+ workflowActionService.addOperator(mockScanPredicate, mockPoint);
+ workflowActionService.addOperator(mockResultPredicate, mockPoint);
const changeStatistics = vi.spyOn(jointUIService,
"changeOperatorStatistics");
- (workflowStatusService as any).statusSubject.next({
- [mockScanPredicate.operatorID]: statisticsIn(OperatorState.Running),
- [mockResultPredicate.operatorID]:
statisticsIn(OperatorState.Completed),
+ (workflowStatusService as any).statisticsSubject.next({
+ [mockScanPredicate.operatorID]: emptyStatistics,
});
- expect(changeStatistics).toHaveBeenCalledWith(
- component.paper,
- mockScanPredicate.operatorID,
- statisticsIn(OperatorState.Running),
- true,
- false
- );
- expect(changeStatistics).toHaveBeenCalledWith(
+ expect(changeStatistics).toHaveBeenCalledWith(component.paper,
mockScanPredicate.operatorID, emptyStatistics);
+ // an operator missing from the payload is forwarded as undefined (a
no-op render)
+ expect(changeStatistics).toHaveBeenCalledWith(component.paper,
mockResultPredicate.operatorID, undefined);
+ });
+
+ it("paints each operator's reported state, defaulting missing ones to
Uninitialized", () => {
+ workflowActionService.addOperator(mockScanPredicate, mockPoint);
+ workflowActionService.addOperator(mockResultPredicate, mockPoint);
+ const changeState = vi.spyOn(jointUIService, "changeOperatorState");
+
+ (workflowStatusService as any).stateSubject.next({
+ [mockScanPredicate.operatorID]: OperatorState.Running,
+ });
+
+ expect(changeState).toHaveBeenCalledWith(component.paper,
mockScanPredicate.operatorID, OperatorState.Running);
+ expect(changeState).toHaveBeenCalledWith(
component.paper,
mockResultPredicate.operatorID,
- statisticsIn(OperatorState.Completed),
- false,
- true
+ OperatorState.Uninitialized
);
});
@@ -1954,41 +1992,33 @@ describe("WorkflowEditorComponent editor wiring", () =>
{
vi.spyOn(executeWorkflowService, "getExecutionState").mockReturnValue({
state: ExecutionState.Recovering,
} as any);
- const changeStatistics = vi.spyOn(jointUIService,
"changeOperatorStatistics");
+ const changeState = vi.spyOn(jointUIService, "changeOperatorState");
- (workflowStatusService as any).statusSubject.next({
- [mockScanPredicate.operatorID]: statisticsIn(OperatorState.Running),
+ (workflowStatusService as any).stateSubject.next({
+ [mockScanPredicate.operatorID]: OperatorState.Running,
});
- expect(changeStatistics).toHaveBeenCalledWith(
- component.paper,
- mockScanPredicate.operatorID,
- expect.objectContaining({ operatorState: OperatorState.Recovering }),
- true,
- false
- );
+ expect(changeState).toHaveBeenCalledWith(component.paper,
mockScanPredicate.operatorID, OperatorState.Recovering);
});
- it("does not invent statistics for an operator missing from the status
payload", () => {
- // The isDefined guard matters most while recovering: without it the
operator would be
- // handed a synthesized `{ operatorState: Recovering }` instead of
nothing at all.
+ it("does not mask an operator missing from the state payload to
Recovering", () => {
+ // Recovering overrides a *reported* state; an operator absent from the
payload
+ // keeps the plain Uninitialized fallback even while the execution
recovers.
workflowActionService.addOperator(mockScanPredicate, mockPoint);
workflowActionService.addOperator(mockResultPredicate, mockPoint);
vi.spyOn(executeWorkflowService, "getExecutionState").mockReturnValue({
state: ExecutionState.Recovering,
} as any);
- const changeStatistics = vi.spyOn(jointUIService,
"changeOperatorStatistics");
+ const changeState = vi.spyOn(jointUIService, "changeOperatorState");
- (workflowStatusService as any).statusSubject.next({
- [mockScanPredicate.operatorID]: statisticsIn(OperatorState.Running),
+ (workflowStatusService as any).stateSubject.next({
+ [mockScanPredicate.operatorID]: OperatorState.Running,
});
- expect(changeStatistics).toHaveBeenCalledWith(
+ expect(changeState).toHaveBeenCalledWith(
component.paper,
mockResultPredicate.operatorID,
- undefined,
- false,
- true
+ OperatorState.Uninitialized
);
});
diff --git
a/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.ts
b/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.ts
index 9e38d8a782..9226476256 100644
---
a/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.ts
+++
b/frontend/src/app/workspace/component/workflow-editor/workflow-editor.component.ts
@@ -46,7 +46,6 @@ import { NzContextMenuService, NzDropdownMenuComponent } from
"ng-zorro-antd/dro
import { ActivatedRoute, Router } from "@angular/router";
import * as _ from "lodash";
import * as joint from "jointjs";
-import { isDefined } from "../../../common/util/predicate";
import { GuiConfigService } from "../../../common/service/gui-config.service";
import { line, curveCatmullRomClosed } from "d3-shape";
import concaveman from "concaveman";
@@ -232,7 +231,7 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
this.handleOperatorSelectionEvents();
this.handlePortHighlightEvent();
this.registerPortDisplayNameChangeHandler();
- this.handleOperatorStatisticsUpdate();
+ this.handleOperatorStatusUpdate();
this.handleHeatmapOverlay();
this.handleHeatmapHover();
this.handleRegionEvents();
@@ -344,47 +343,41 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
}
/**
- * This method subscribe to workflowStatusService's status stream
- * for Each processStatus that has been emitted
- * 1. enable operatorStatusTooltipDisplay because tooltip will not be
empty
- * 2. for each operator in current texeraGraph:
- * - find its Statistics in processStatus, thrown an error if not
found
- * - generate its corresponding tooltip's id
- * - pass the tooltip id and Statistics to jointUIService
- * the specific tooltip content will be updated
- * - if operator is in a group, save statistics in group's
operatorInfo
- * 3. Whenever a group is expanded
- * - for each operatorInfo, display statistics if there are some
saved.
+ * Renders WorkflowStatusService's two separate sub-concepts on the paper:
+ * - the state stream drives each operator's execution-state rendering
+ * (see {@link effectiveOperatorState} for the fallback/override rules)
+ * - the statistics stream drives each operator's port row-count labels
+ * and worker count.
*/
- private handleOperatorStatisticsUpdate(): void {
+ private handleOperatorStatusUpdate(): void {
this.workflowStatusService
- .getStatusUpdateStream()
+ .getStateUpdateStream()
.pipe(untilDestroyed(this))
- .subscribe(status => {
+ .subscribe(state => {
this.workflowActionService
.getTexeraGraph()
- .getAllOperators()
- .forEach(op => {
- if (
- isDefined(status[op.operatorID]) &&
- this.executeWorkflowService.getExecutionState().state ===
ExecutionState.Recovering
- ) {
- status[op.operatorID] = {
- ...status[op.operatorID],
- operatorState: OperatorState.Recovering,
- };
- }
-
- this.jointUIService.changeOperatorStatistics(
+ .getAllOperatorIDs()
+ .forEach(operatorID => {
+ this.jointUIService.changeOperatorState(
this.paper,
- op.operatorID,
- status[op.operatorID],
- this.isSource(op.operatorID),
- this.isSink(op.operatorID)
+ operatorID,
+ this.effectiveOperatorState(state[operatorID])
);
});
});
+ this.workflowStatusService
+ .getStatisticsUpdateStream()
+ .pipe(untilDestroyed(this))
+ .subscribe(statistics => {
+ this.workflowActionService
+ .getTexeraGraph()
+ .getAllOperatorIDs()
+ .forEach(operatorID => {
+ this.jointUIService.changeOperatorStatistics(this.paper,
operatorID, statistics[operatorID]);
+ });
+ });
+
this.executeWorkflowService
.getExecutionStateStream()
.pipe(untilDestroyed(this))
@@ -402,9 +395,9 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
}
this.workflowActionService
.getTexeraGraph()
- .getAllOperators()
- .forEach(op => {
- this.jointUIService.changeOperatorState(this.paper,
op.operatorID, operatorState);
+ .getAllOperatorIDs()
+ .forEach(operatorID => {
+ this.jointUIService.changeOperatorState(this.paper, operatorID,
operatorState);
});
}
});
@@ -414,24 +407,19 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
// operators are recreated from the workflow JSON — restore their visual
// state from the cached status so completed runs don't appear to reset.
// Restores port labels / worker count via changeOperatorStatistics, then
- // delegates the final border color to applyOperatorBorder so the same
- // priority rules apply as for the validation pass.
+ // delegates the execution-state rendering and the final border color to
+ // applyOperatorStateAndBorder so the same priority rules apply as for the
+ // validation pass.
this.workflowActionService
.getTexeraGraph()
.getOperatorAddStream()
.pipe(untilDestroyed(this))
.subscribe(operator => {
- const statistics =
this.workflowStatusService.getCurrentStatus()[operator.operatorID];
+ const statistics =
this.workflowStatusService.getCurrentStatistics()[operator.operatorID];
if (statistics) {
- this.jointUIService.changeOperatorStatistics(
- this.paper,
- operator.operatorID,
- statistics,
- this.isSource(operator.operatorID),
- this.isSink(operator.operatorID)
- );
+ this.jointUIService.changeOperatorStatistics(this.paper,
operator.operatorID, statistics);
}
- this.applyOperatorBorder(
+ this.applyOperatorStateAndBorder(
operator.operatorID,
this.validationWorkflowService.validateOperator(operator.operatorID)
);
@@ -576,16 +564,34 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
}
/**
- * Single source of truth for the operator's border color. Both the
- * validation stream and the operator-add stream route through here so
- * the priority order is consistent regardless of which event fires last:
- * 1. Invalid operator → red (validation takes priority).
- * 2. Valid operator with a cached execution status → execution-state
color.
- * 3. Valid operator with no cached status → default valid (gray).
+ * Resolves the execution state to render for an operator: an operator the
+ * state map does not know falls back to Uninitialized, and a reported state
+ * is masked to Recovering while the workflow execution is recovering.
+ */
+ private effectiveOperatorState(reportedState?: OperatorState): OperatorState
{
+ if (reportedState === undefined) {
+ return OperatorState.Uninitialized;
+ }
+ if (this.executeWorkflowService.getExecutionState().state ===
ExecutionState.Recovering) {
+ return OperatorState.Recovering;
+ }
+ return reportedState;
+ }
+
+ /**
+ * Single source of truth for the operator's state rendering and border
+ * color. Both the validation stream and the operator-add stream route
+ * through here so the priority order is consistent regardless of which
+ * event fires last:
+ * 1. A cached execution state repaints the full state rendering first
+ * (state label, fills), so it survives navigating away and back.
+ * 2. Invalid operator → red stroke on top (validation takes priority).
+ * 3. Valid operator with a cached execution state keeps its state stroke.
+ * 4. Valid operator with no cached state → default valid (gray).
*
* Centralizing this here avoids the race where the validation pass
* overwrites a state-derived stroke (or vice versa) for an operator that
- * is both invalid and has a cached execution status.
+ * is both invalid and has a cached execution state.
*
* Both callers obtain the Validation themselves and pass it in: the
* validation-stream subscriber forwards the result the stream just emitted,
@@ -593,15 +599,16 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
* the parameter required means the color decision never silently depends on
* a recompute hidden inside this helper.
*/
- private applyOperatorBorder(operatorID: string, validation: Validation):
void {
+ private applyOperatorStateAndBorder(operatorID: string, validation:
Validation): void {
+ const operatorState =
this.workflowStatusService.getCurrentState()[operatorID];
+ if (operatorState) {
+ this.jointUIService.changeOperatorState(this.paper, operatorID,
this.effectiveOperatorState(operatorState));
+ }
if (!validation.isValid) {
this.jointUIService.changeOperatorColor(this.paper, operatorID, false);
return;
}
- const statistics =
this.workflowStatusService.getCurrentStatus()[operatorID];
- if (statistics) {
- this.jointUIService.changeOperatorState(this.paper, operatorID,
statistics.operatorState);
- } else {
+ if (!operatorState) {
this.jointUIService.changeOperatorColor(this.paper, operatorID, true);
}
}
@@ -1259,14 +1266,14 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
/**
* Applies the validation result to the operator's border. Delegates to
- * applyOperatorBorder so validation, cached-execution-status, and the
+ * applyOperatorStateAndBorder so validation, cached-execution-state, and the
* default-valid case are decided in one place.
*/
private handleOperatorValidation(): void {
this.validationWorkflowService
.getOperatorValidationStream()
.pipe(untilDestroyed(this))
- .subscribe(value => this.applyOperatorBorder(value.operatorID,
value.validation));
+ .subscribe(value => this.applyOperatorStateAndBorder(value.operatorID,
value.validation));
// Operators already in the graph when this editor mounts produced no add
event and won't hit
// the validation stream until they change, so nothing corrects the border
they were drawn
@@ -1285,22 +1292,13 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
private paintCurrentOperatorState(): void {
this.workflowActionService
.getTexeraGraph()
- .getAllOperators()
- .forEach(operator => {
- const statistics =
this.workflowStatusService.getCurrentStatus()[operator.operatorID];
+ .getAllOperatorIDs()
+ .forEach(operatorID => {
+ const statistics =
this.workflowStatusService.getCurrentStatistics()[operatorID];
if (statistics) {
- this.jointUIService.changeOperatorStatistics(
- this.paper,
- operator.operatorID,
- statistics,
- this.isSource(operator.operatorID),
- this.isSink(operator.operatorID)
- );
+ this.jointUIService.changeOperatorStatistics(this.paper, operatorID,
statistics);
}
- this.applyOperatorBorder(
- operator.operatorID,
- this.validationWorkflowService.validateOperator(operator.operatorID)
- );
+ this.applyOperatorStateAndBorder(operatorID,
this.validationWorkflowService.validateOperator(operatorID));
});
}
/* v8 ignore stop */
@@ -1690,14 +1688,6 @@ export class WorkflowEditorComponent implements OnInit,
AfterViewInit, OnDestroy
});
}
- private isSource(operatorID: string): boolean {
- return
this.workflowActionService.getTexeraGraph().getOperator(operatorID).inputPorts.length
== 0;
- }
-
- private isSink(operatorID: string): boolean {
- return
this.workflowActionService.getTexeraGraph().getOperator(operatorID).outputPorts.length
== 0;
- }
-
/**
* Handles mouse events to enable shared cursor.
*/
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 288846346d..6a7a4806b2 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
@@ -743,15 +743,29 @@ describe("JointUIService", () => {
return { paper, attrSpy, portPropSpy };
}
- it("falls back to the Uninitialized state when statistics is undefined",
() => {
+ it("renders nothing when statistics is undefined (state is rendered
separately)", () => {
+ const { paper, attrSpy, portPropSpy } = makeStatsPaper(() => []);
+ const service = new JointUIService(emptyMetadataStub as never);
+ service.changeOperatorStatistics(paper, "op-1", undefined);
+ expect(attrSpy).not.toHaveBeenCalled();
+ expect(portPropSpy).not.toHaveBeenCalled();
+ });
+
+ it("does not touch the operator's execution-state rendering", () => {
+ // State is a separate sub-concept rendered via changeOperatorState;
+ // a statistics update must not repaint the state class.
const { paper, attrSpy } = makeStatsPaper(() => []);
const service = new JointUIService(emptyMetadataStub as never);
- service.changeOperatorStatistics(paper, "op-1", undefined, false, false);
- // changeOperatorState writes the state-class fill payload.
- const [payload] = attrSpy.mock.calls[0];
- expect((payload as Record<string, { text: string
}>)[`.${operatorStateClass}`].text).toBe(
- OperatorState.Uninitialized.toString()
+ service.changeOperatorStatistics(paper, "op-1", {
+ aggregatedInputRowCount: 0,
+ aggregatedOutputRowCount: 0,
+ inputPortMetrics: {},
+ outputPortMetrics: {},
+ });
+ const stateWrites = attrSpy.mock.calls.filter(
+ c => typeof c[0] === "object" && c[0] !== null &&
`.${operatorStateClass}` in (c[0] as object)
);
+ expect(stateWrites).toHaveLength(0);
});
it("writes per-port counts derived from inputPortMetrics and
outputPortMetrics", () => {
@@ -762,20 +776,13 @@ describe("JointUIService", () => {
{ id: "out-1", group: "out", attrs: { ".port-label": { text: "result:
0" } } },
]);
const service = new JointUIService(emptyMetadataStub as never);
- service.changeOperatorStatistics(
- paper,
- "op-1",
- {
- operatorState: OperatorState.Running,
- aggregatedInputRowCount: 0,
- aggregatedOutputRowCount: 0,
- inputPortMetrics: { "0": 42 },
- outputPortMetrics: { "1": 7 },
- numWorkers: 3,
- },
- false,
- false
- );
+ service.changeOperatorStatistics(paper, "op-1", {
+ aggregatedInputRowCount: 0,
+ aggregatedOutputRowCount: 0,
+ inputPortMetrics: { "0": 42 },
+ outputPortMetrics: { "1": 7 },
+ numWorkers: 3,
+ });
expect(portPropSpy).toHaveBeenCalledWith("in-0",
"attrs/.port-label/text", (42).toLocaleString());
expect(portPropSpy).toHaveBeenCalledWith("out-1",
"attrs/.port-label/text", (7).toLocaleString());
});
@@ -783,20 +790,13 @@ describe("JointUIService", () => {
it("writes the worker count label when statistics include numWorkers", ()
=> {
const { paper, attrSpy } = makeStatsPaper(() => []);
const service = new JointUIService(emptyMetadataStub as never);
- service.changeOperatorStatistics(
- paper,
- "op-1",
- {
- operatorState: OperatorState.Ready,
- aggregatedInputRowCount: 0,
- aggregatedOutputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- numWorkers: 8,
- },
- false,
- false
- );
+ service.changeOperatorStatistics(paper, "op-1", {
+ aggregatedInputRowCount: 0,
+ aggregatedOutputRowCount: 0,
+ inputPortMetrics: {},
+ outputPortMetrics: {},
+ numWorkers: 8,
+ });
// attr() is called once with the workers selector and the formatted
string.
const valuesWritten = attrSpy.mock.calls.map(c => c[1]);
expect(valuesWritten).toContain("#workers: 8");
@@ -805,19 +805,12 @@ describe("JointUIService", () => {
it("defaults the worker count to 1 when numWorkers is unspecified", () => {
const { paper, attrSpy } = makeStatsPaper(() => []);
const service = new JointUIService(emptyMetadataStub as never);
- service.changeOperatorStatistics(
- paper,
- "op-1",
- {
- operatorState: OperatorState.Ready,
- aggregatedInputRowCount: 0,
- aggregatedOutputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- false,
- false
- );
+ service.changeOperatorStatistics(paper, "op-1", {
+ aggregatedInputRowCount: 0,
+ aggregatedOutputRowCount: 0,
+ inputPortMetrics: {},
+ outputPortMetrics: {},
+ });
const valuesWritten = attrSpy.mock.calls.map(c => c[1]);
expect(valuesWritten).toContain("#workers: 1");
});
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 29ddb1f54a..bb3ba9e57f 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
@@ -367,20 +367,20 @@ export class JointUIService {
return operatorElement;
}
+ /**
+ * Renders the statistics sub-concept only (port row counts and worker
+ * count); the operator's execution state is rendered separately via
+ * {@link changeOperatorState}.
+ */
public changeOperatorStatistics(
jointPaper: joint.dia.Paper,
operatorID: string,
- statistics: OperatorStatistics | undefined,
- isSource: boolean,
- isSink: boolean
+ statistics: OperatorStatistics | undefined
): void {
if (!statistics) {
- this.changeOperatorState(jointPaper, operatorID,
OperatorState.Uninitialized);
return;
}
- this.changeOperatorState(jointPaper, operatorID, statistics.operatorState);
-
const element = jointPaper.getModelById(operatorID) as
joint.shapes.devs.Model;
const allPorts = element.getPorts();
const inPorts = allPorts.filter(p => p.group === "in");
@@ -411,7 +411,6 @@ export class JointUIService {
element.portProp(portId, "attrs/.port-label/text",
count.toLocaleString());
}
});
- this.changeOperatorState(jointPaper, operatorID, statistics.operatorState);
}
public foldOperatorDetails(jointPaper: joint.dia.Paper, operatorID: string):
void {
jointPaper.getModelById(operatorID).attr({
diff --git
a/frontend/src/app/workspace/service/operator-debug/udf-debug.service.spec.ts
b/frontend/src/app/workspace/service/operator-debug/udf-debug.service.spec.ts
index 5f22bd3dbd..47b6097479 100644
---
a/frontend/src/app/workspace/service/operator-debug/udf-debug.service.spec.ts
+++
b/frontend/src/app/workspace/service/operator-debug/udf-debug.service.spec.ts
@@ -24,7 +24,7 @@ import { WorkflowActionService } from
"../workflow-graph/model/workflow-action.s
import { WorkflowStatusService } from
"../workflow-status/workflow-status.service";
import { ExecuteWorkflowService } from
"../execute-workflow/execute-workflow.service";
import { Observable, Subject } from "rxjs";
-import { OperatorState, OperatorStatistics } from
"../../types/execute-workflow.interface";
+import { OperatorState } from "../../types/execute-workflow.interface";
import { WorkflowGraphReadonly } from "../workflow-graph/model/workflow-graph";
import { mockPoint, mockPythonUDFPredicate } from
"../workflow-graph/model/mock-workflow-data";
import { OperatorMetadataService } from
"../operator-metadata/operator-metadata.service";
@@ -40,7 +40,7 @@ describe("UdfDebugServiceSpec", () => {
let mockWorkflowWebsocketService: Mocked<WorkflowWebsocketService>;
let mockWorkflowStatusService: Mocked<WorkflowStatusService>;
let mockExecuteWorkflowService: Mocked<ExecuteWorkflowService>;
- let statusUpdateStream: Subject<Record<string, OperatorStatistics>>;
+ let stateUpdateStream: Subject<Record<string, OperatorState>>;
let consoleUpdateEventStream: Subject<ConsoleUpdateEvent>;
let texeraGraph: WorkflowGraphReadonly;
let stubWorker = "worker1";
@@ -51,15 +51,15 @@ describe("UdfDebugServiceSpec", () => {
send: vi.fn(),
subscribeToEvent: vi.fn(),
} as unknown as Mocked<WorkflowWebsocketService>;
- mockWorkflowStatusService = { getStatusUpdateStream: vi.fn() } as unknown
as Mocked<WorkflowStatusService>;
+ mockWorkflowStatusService = { getStateUpdateStream: vi.fn() } as unknown
as Mocked<WorkflowStatusService>;
mockExecuteWorkflowService = { getWorkerIds: vi.fn() } as unknown as
Mocked<ExecuteWorkflowService>;
// Initialize the mock streams
- statusUpdateStream = new Subject();
+ stateUpdateStream = new Subject();
consoleUpdateEventStream = new Subject();
// Set mock return values
-
mockWorkflowStatusService.getStatusUpdateStream.mockReturnValue(statusUpdateStream.asObservable());
+
mockWorkflowStatusService.getStateUpdateStream.mockReturnValue(stateUpdateStream.asObservable());
mockWorkflowWebsocketService.subscribeToEvent.mockReturnValue(
consoleUpdateEventStream.asObservable() as
Observable<TexeraWebsocketEvent>
);
@@ -93,7 +93,7 @@ describe("UdfDebugServiceSpec", () => {
afterEach(() => {
// Clean up the streams after each test
- statusUpdateStream.complete();
+ stateUpdateStream.complete();
consoleUpdateEventStream.complete();
});
@@ -198,15 +198,7 @@ describe("UdfDebugServiceSpec", () => {
const debugState =
service.getDebugState(mockPythonUDFPredicate.operatorID);
const operatorId = mockPythonUDFPredicate.operatorID;
debugState.set(operatorId, { breakpointId: 1, condition: "x > 5", hit:
false });
- statusUpdateStream.next({
- [operatorId]: {
- operatorState: OperatorState.Uninitialized,
- aggregatedInputRowCount: 0,
- aggregatedOutputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- },
- });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Uninitialized });
expect(debugState.size).toBe(0);
});
@@ -536,15 +528,8 @@ describe("UdfDebugServiceSpec", () => {
const debugState = service.getDebugState(operatorId);
debugState.set("10", { breakpointId: 1, condition: "", hit: false });
- const running: OperatorStatistics = {
- operatorState: OperatorState.Running,
- aggregatedInputRowCount: 0,
- aggregatedOutputRowCount: 0,
- inputPortMetrics: {},
- outputPortMetrics: {},
- };
- statusUpdateStream.next({ [operatorId]: running });
- statusUpdateStream.next({ "some-other-operator": { ...running,
operatorState: OperatorState.Uninitialized } });
+ stateUpdateStream.next({ [operatorId]: OperatorState.Running });
+ stateUpdateStream.next({ "some-other-operator":
OperatorState.Uninitialized });
expect(debugState.size).toBe(1);
});
diff --git
a/frontend/src/app/workspace/service/operator-debug/udf-debug.service.ts
b/frontend/src/app/workspace/service/operator-debug/udf-debug.service.ts
index 4d931b2780..1e19fb54bb 100644
--- a/frontend/src/app/workspace/service/operator-debug/udf-debug.service.ts
+++ b/frontend/src/app/workspace/service/operator-debug/udf-debug.service.ts
@@ -155,8 +155,8 @@ export class UdfDebugService {
*/
private registerOperatorStateChangeHandler(operatorId: string) {
this.workflowStatusService
- .getStatusUpdateStream()
- .pipe(filter(event => event[operatorId]?.operatorState ===
OperatorState.Uninitialized))
+ .getStateUpdateStream()
+ .pipe(filter(event => event[operatorId] === OperatorState.Uninitialized))
.subscribe(() => this.getDebugState(operatorId).clear());
}
diff --git
a/frontend/src/app/workspace/service/workflow-graph/model/workflow-graph.ts
b/frontend/src/app/workspace/service/workflow-graph/model/workflow-graph.ts
index 4d4bb93ede..9c77ad4170 100644
--- a/frontend/src/app/workspace/service/workflow-graph/model/workflow-graph.ts
+++ b/frontend/src/app/workspace/service/workflow-graph/model/workflow-graph.ts
@@ -634,6 +634,15 @@ export class WorkflowGraph {
);
}
+ /**
+ * Returns the IDs of all operators in the graph. Unlike {@link
getAllOperators},
+ * this does not materialize the operator predicates from the shared model,
so
+ * it is cheap enough for per-update paths that only need the IDs.
+ */
+ public getAllOperatorIDs(): readonly string[] {
+ return Array.from(this.sharedModel.operatorIDMap.keys() as
IterableIterator<string>);
+ }
+
/**
* Returns an array of all enabled operators in the graph.
*/
diff --git
a/frontend/src/app/workspace/service/workflow-status/performance-metrics.spec.ts
b/frontend/src/app/workspace/service/workflow-status/performance-metrics.spec.ts
index 47c018e9d9..3ef562e851 100644
---
a/frontend/src/app/workspace/service/workflow-status/performance-metrics.spec.ts
+++
b/frontend/src/app/workspace/service/workflow-status/performance-metrics.spec.ts
@@ -18,14 +18,13 @@
*/
import { OperatorPerformanceMetrics, extractPerformanceMetrics } from
"./performance-metrics";
-import { OperatorState, OperatorStatistics } from
"../../types/execute-workflow.interface";
+import { OperatorStatistics } from "../../types/execute-workflow.interface";
/**
* A complete statistics object, mirroring what the backend sends once every
* field is typed. All five timing/size fields are present.
*/
const fullStats: OperatorStatistics = {
- operatorState: OperatorState.Running,
aggregatedInputRowCount: 1_000_000,
aggregatedInputSize: 84_000_000,
inputPortMetrics: { "0": 1_000_000 },
@@ -44,7 +43,6 @@ const fullStats: OperatorStatistics = {
* The mapper must survive this without emitting NaN/undefined.
*/
const partialStats: OperatorStatistics = {
- operatorState: OperatorState.Uninitialized,
aggregatedInputRowCount: 0,
inputPortMetrics: {},
aggregatedOutputRowCount: 0,
diff --git
a/frontend/src/app/workspace/service/workflow-status/performance-metrics.ts
b/frontend/src/app/workspace/service/workflow-status/performance-metrics.ts
index 8f2b4ffa46..6c87952f5b 100644
--- a/frontend/src/app/workspace/service/workflow-status/performance-metrics.ts
+++ b/frontend/src/app/workspace/service/workflow-status/performance-metrics.ts
@@ -23,9 +23,10 @@ import { OperatorStatistics } from
"../../types/execute-workflow.interface";
* Derived per-operator performance metrics.
*
* This is the ground-truth model captured by {@link WorkflowStatusService}.
It is
- * a flat projection of the raw {@link OperatorStatistics} the backend streams
over
- * the websocket, with missing optional fields defaulted. Keyed by operator id
at
- * the map level (mirroring {@link OperatorStatistics}), so the id is not
repeated here.
+ * a flat projection of the {@link OperatorStatistics} sub-concept the service
+ * splits out of the OperatorRuntimeStatus objects the backend streams over the
+ * websocket, with missing optional fields defaulted. Keyed by operator id at
the
+ * map level (mirroring {@link OperatorStatistics}), so the id is not repeated
here.
*/
export interface OperatorPerformanceMetrics
extends Readonly<{
@@ -40,7 +41,7 @@ export interface OperatorPerformanceMetrics
}> {}
/**
- * Project a single raw {@link OperatorStatistics} into the flat performance
model.
+ * Project a single {@link OperatorStatistics} into the flat performance model.
*
* Several fields are optional on {@link OperatorStatistics} because the
frontend
* builds partial objects (e.g. WorkflowStatusService.resetStatus, which omits
the
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 a09f1b41f7..7aa5b6aa53 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
@@ -22,11 +22,10 @@ import { Subject } from "rxjs";
import { WorkflowStatusService } from "./workflow-status.service";
import { WorkflowWebsocketService } from
"../workflow-websocket/workflow-websocket.service";
import { OperatorPerformanceMetrics } from "./performance-metrics";
-import { OperatorState, OperatorStatistics } from
"../../types/execute-workflow.interface";
+import { OperatorRuntimeStatus, OperatorState, OperatorStatistics } from
"../../types/execute-workflow.interface";
import { TexeraWebsocketEvent } from
"../../types/workflow-websocket.interface";
-const sampleStats: OperatorStatistics = {
- operatorState: OperatorState.Running,
+const sampleStatistics: OperatorStatistics = {
aggregatedInputRowCount: 1_000,
aggregatedInputSize: 8_000,
inputPortMetrics: { "0": 1_000 },
@@ -39,7 +38,13 @@ const sampleStats: OperatorStatistics = {
aggregatedIdleTime: 700_000,
};
-function statsEvent(operatorStatistics: Record<string, OperatorStatistics>):
TexeraWebsocketEvent {
+// The wire object the engine streams: state and statistics bundled together.
+const sampleRuntimeStatus: OperatorRuntimeStatus = {
+ operatorState: OperatorState.Running,
+ ...sampleStatistics,
+};
+
+function statsEvent(operatorStatistics: Record<string,
OperatorRuntimeStatus>): TexeraWebsocketEvent {
return { type: "OperatorStatisticsUpdateEvent", operatorStatistics } as
TexeraWebsocketEvent;
}
@@ -58,32 +63,82 @@ describe("WorkflowStatusService", () => {
service = TestBed.inject(WorkflowStatusService);
});
- it("forwards an OperatorStatisticsUpdateEvent to the status stream", () => {
- const received: Record<string, OperatorStatistics>[] = [];
- service.getStatusUpdateStream().subscribe(s => received.push(s));
+ it("splits an OperatorStatisticsUpdateEvent into the state and statistics
streams", () => {
+ const stateEmissions: Record<string, OperatorState>[] = [];
+ const statisticsEmissions: Record<string, OperatorStatistics>[] = [];
+ service.getStateUpdateStream().subscribe(s => stateEmissions.push(s));
+ service.getStatisticsUpdateStream().subscribe(s =>
statisticsEmissions.push(s));
- websocketEventSubject.next(statsEvent({ op1: sampleStats }));
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
- expect(received).toHaveLength(1);
- expect(received[0]).toEqual({ op1: sampleStats });
- expect(service.getCurrentStatus()).toEqual({ op1: sampleStats });
+ expect(stateEmissions).toHaveLength(1);
+ expect(stateEmissions[0]).toEqual({ op1: OperatorState.Running });
+ expect(statisticsEmissions).toHaveLength(1);
+ expect(statisticsEmissions[0]).toEqual({ op1: sampleStatistics });
+ expect(service.getCurrentState()).toEqual({ op1: OperatorState.Running });
+ expect(service.getCurrentStatistics()).toEqual({ op1: sampleStatistics });
+ });
+
+ it("does not leak the operator state into the statistics concept", () => {
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
+
expect(service.getCurrentStatistics()["op1"]).not.toHaveProperty("operatorState");
+ });
+
+ it("exposes state and statistics as independent streams", () => {
+ // A consumer subscribed to only one of the two sub-concepts sees exactly
+ // one emission per update, unaffected by the other stream.
+ let stateCount = 0;
+ let statisticsCount = 0;
+ service.getStateUpdateStream().subscribe(() => stateCount++);
+ service.getStatisticsUpdateStream().subscribe(() => statisticsCount++);
+
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
+ websocketEventSubject.next(statsEvent({ op1: { ...sampleRuntimeStatus,
operatorState: OperatorState.Paused } }));
+
+ expect(stateCount).toBe(2);
+ expect(statisticsCount).toBe(2);
+ expect(service.getCurrentState()).toEqual({ op1: OperatorState.Paused });
+ });
+
+ it("emits state before statistics, so a statistics subscriber sees the
matching state snapshot", () => {
+ const order: string[] = [];
+ // Captured in the subscriber, asserted after next() returns: rxjs
re-throws
+ // a subscriber's error asynchronously, so an expect() inside the callback
+ // could not fail this test.
+ let stateSeenByStatisticsSubscriber: Record<string, OperatorState> |
undefined;
+ service.getStateUpdateStream().subscribe(() => order.push("state"));
+ service.getStatisticsUpdateStream().subscribe(() => {
+ order.push("statistics");
+ stateSeenByStatisticsSubscriber = service.getCurrentState();
+ });
+
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
+
+ expect(order).toEqual(["state", "statistics"]);
+ // The licensed read: by the time statistics arrive, the state snapshot
+ // already reflects the same wire event. The converse does not hold.
+ expect(stateSeenByStatisticsSubscriber).toEqual({ op1:
OperatorState.Running });
});
it("ignores websocket events of other types", () => {
- const received: Record<string, OperatorStatistics>[] = [];
- service.getStatusUpdateStream().subscribe(s => received.push(s));
+ const stateEmissions: Record<string, OperatorState>[] = [];
+ const statisticsEmissions: Record<string, OperatorStatistics>[] = [];
+ service.getStateUpdateStream().subscribe(s => stateEmissions.push(s));
+ service.getStatisticsUpdateStream().subscribe(s =>
statisticsEmissions.push(s));
websocketEventSubject.next({ type: "WorkflowErrorEvent" } as unknown as
TexeraWebsocketEvent);
- expect(received).toHaveLength(0);
- expect(service.getCurrentStatus()).toEqual({});
+ expect(stateEmissions).toHaveLength(0);
+ expect(statisticsEmissions).toHaveLength(0);
+ expect(service.getCurrentState()).toEqual({});
+ expect(service.getCurrentStatistics()).toEqual({});
});
- it("derives performance metrics from a status update", () => {
+ it("derives performance metrics from a statistics update", () => {
const emissions: Record<string, OperatorPerformanceMetrics>[] = [];
service.getPerformanceMetricsStream().subscribe(m => emissions.push(m));
- websocketEventSubject.next(statsEvent({ op1: sampleStats }));
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
// BehaviorSubject seeds {} then emits the derived map.
const latest = emissions[emissions.length - 1];
@@ -102,7 +157,7 @@ describe("WorkflowStatusService", () => {
it("keys the derived metrics by operator id, including unicode ids", () => {
const id = "算子-✓-1";
- websocketEventSubject.next(statsEvent({ [id]: sampleStats }));
+ websocketEventSubject.next(statsEvent({ [id]: sampleRuntimeStatus }));
expect(Object.keys(service.getCurrentPerformanceMetrics())).toEqual([id]);
});
@@ -114,7 +169,7 @@ describe("WorkflowStatusService", () => {
});
it("defaults missing optional fields to 0 when deriving metrics", () => {
- const partial: OperatorStatistics = {
+ const partial: OperatorRuntimeStatus = {
operatorState: OperatorState.Uninitialized,
aggregatedInputRowCount: 0,
inputPortMetrics: {},
@@ -132,12 +187,20 @@ describe("WorkflowStatusService", () => {
expect(m.numWorkers).toBe(1);
});
- it("resetStatus zeros the metrics for known operators", () => {
- websocketEventSubject.next(statsEvent({ op1: sampleStats }));
+ it("resetStatus resets both concepts for known operators", () => {
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
service.resetStatus();
- const m = service.getCurrentPerformanceMetrics()["op1"];
- expect(m).toEqual({
+ expect(service.getCurrentState()).toEqual({ op1:
OperatorState.Uninitialized });
+ expect(service.getCurrentStatistics()).toEqual({
+ op1: {
+ aggregatedInputRowCount: 0,
+ inputPortMetrics: {},
+ aggregatedOutputRowCount: 0,
+ outputPortMetrics: {},
+ },
+ });
+ expect(service.getCurrentPerformanceMetrics()["op1"]).toEqual({
dataProcessingTimeNs: 0,
controlProcessingTimeNs: 0,
idleTimeNs: 0,
@@ -149,11 +212,12 @@ describe("WorkflowStatusService", () => {
});
});
- it("clearStatus empties both the status and performance-metrics snapshots",
() => {
- websocketEventSubject.next(statsEvent({ op1: sampleStats }));
+ it("clearStatus empties the state, statistics, and performance-metrics
snapshots", () => {
+ websocketEventSubject.next(statsEvent({ op1: sampleRuntimeStatus }));
service.clearStatus();
- expect(service.getCurrentStatus()).toEqual({});
+ expect(service.getCurrentState()).toEqual({});
+ expect(service.getCurrentStatistics()).toEqual({});
expect(service.getCurrentPerformanceMetrics()).toEqual({});
});
});
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 709768deff..729adb6507 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
@@ -27,9 +27,15 @@ import { OperatorPerformanceMetrics,
extractPerformanceMetrics } from "./perform
providedIn: "root",
})
export class WorkflowStatusService {
- // status is responsible for passing websocket responses to other components
- private statusSubject = new Subject<Record<string, OperatorStatistics>>();
- private currentStatus: Record<string, OperatorStatistics> = {};
+ // The engine streams operator state and operator statistics bundled in one
+ // wire object (OperatorRuntimeStatus); this service splits them into two
+ // separate sub-concepts, each with its own stream and snapshot. Derived
+ // performance metrics are the third, separate concept.
+ private stateSubject = new Subject<Record<string, OperatorState>>();
+ private currentState: Record<string, OperatorState> = {};
+
+ private statisticsSubject = new Subject<Record<string,
OperatorStatistics>>();
+ private currentStatistics: Record<string, OperatorStatistics> = {};
// Derived, ground-truth performance metrics for the heat-map overlay.
Backed by
// a BehaviorSubject so a consumer that subscribes after a run already
streamed
@@ -37,27 +43,57 @@ export class WorkflowStatusService {
private performanceMetricsSubject = new BehaviorSubject<Record<string,
OperatorPerformanceMetrics>>({});
constructor(private workflowWebsocketService: WorkflowWebsocketService) {
- // Single derivation path: every status emission (websocket, reset, clear,
or
- // externally fed historical stats) recomputes the performance metrics.
- this.getStatusUpdateStream().subscribe(status => {
- this.currentStatus = status;
-
this.performanceMetricsSubject.next(this.buildPerformanceMetrics(status));
+ this.getStateUpdateStream().subscribe(state => {
+ this.currentState = state;
+ });
+
+ // Single derivation path: every statistics emission (websocket, reset, or
+ // clear) recomputes the performance metrics.
+ this.getStatisticsUpdateStream().subscribe(statistics => {
+ this.currentStatistics = statistics;
+
this.performanceMetricsSubject.next(this.buildPerformanceMetrics(statistics));
});
+ // Each wire event produces exactly one emission on each stream, state
+ // first and statistics second (resetStatus/clearStatus follow the same
+ // order). The guarantee is one-directional: a statistics subscriber may
+ // read getCurrentState() and see the matching snapshot, but a state
+ // subscriber reading getCurrentStatistics() sees the previous emission.
+ // Pinned by the "emits state before statistics" spec.
this.workflowWebsocketService.websocketEvent().subscribe(event => {
if (event.type !== "OperatorStatisticsUpdateEvent") {
return;
}
- this.statusSubject.next(event.operatorStatistics);
+ 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);
});
}
- public getStatusUpdateStream(): Observable<Record<string,
OperatorStatistics>> {
- return this.statusSubject.asObservable();
+ /** Stream of per-operator execution states, keyed by operator id. */
+ public getStateUpdateStream(): Observable<Record<string, OperatorState>> {
+ return this.stateSubject.asObservable();
+ }
+
+ /** Synchronous snapshot of the latest per-operator execution states. */
+ public getCurrentState(): Record<string, OperatorState> {
+ return this.currentState;
+ }
+
+ /** Stream of per-operator statistics (row counts, sizes, timing), keyed by
operator id. */
+ public getStatisticsUpdateStream(): Observable<Record<string,
OperatorStatistics>> {
+ return this.statisticsSubject.asObservable();
}
- public getCurrentStatus(): Record<string, OperatorStatistics> {
- return this.currentStatus;
+ /** Synchronous snapshot of the latest per-operator statistics. */
+ public getCurrentStatistics(): Record<string, OperatorStatistics> {
+ return this.currentStatistics;
}
/** Stream of derived per-operator performance metrics, keyed by operator
id. */
@@ -71,34 +107,35 @@ export class WorkflowStatusService {
}
private buildPerformanceMetrics(
- status: Record<string, OperatorStatistics>
+ statistics: Record<string, OperatorStatistics>
): Record<string, OperatorPerformanceMetrics> {
const metrics: Record<string, OperatorPerformanceMetrics> = {};
- for (const operatorId of Object.keys(status)) {
- metrics[operatorId] = extractPerformanceMetrics(status[operatorId]);
+ for (const operatorId of Object.keys(statistics)) {
+ metrics[operatorId] = extractPerformanceMetrics(statistics[operatorId]);
}
return metrics;
}
public resetStatus(): void {
- const initStatus: Record<string, OperatorStatistics> =
Object.keys(this.currentStatus).reduce(
- (accumulator, operatorId) => {
- accumulator[operatorId] = {
- operatorState: OperatorState.Uninitialized,
- aggregatedInputRowCount: 0,
- inputPortMetrics: {},
- aggregatedOutputRowCount: 0,
- outputPortMetrics: {},
- };
- return accumulator;
- },
- {} as Record<string, OperatorStatistics>
- );
- this.statusSubject.next(initStatus);
+ const initState: Record<string, OperatorState> = {};
+ for (const operatorId of Object.keys(this.currentState)) {
+ initState[operatorId] = OperatorState.Uninitialized;
+ }
+ const initStatistics: Record<string, OperatorStatistics> = {};
+ for (const operatorId of Object.keys(this.currentStatistics)) {
+ initStatistics[operatorId] = {
+ aggregatedInputRowCount: 0,
+ inputPortMetrics: {},
+ aggregatedOutputRowCount: 0,
+ outputPortMetrics: {},
+ };
+ }
+ this.stateSubject.next(initState);
+ this.statisticsSubject.next(initStatistics);
}
public clearStatus(): void {
- this.currentStatus = {};
- this.statusSubject.next({});
+ this.stateSubject.next({});
+ this.statisticsSubject.next({});
}
}
diff --git a/frontend/src/app/workspace/types/execute-workflow.interface.ts
b/frontend/src/app/workspace/types/execute-workflow.interface.ts
index 8bb7696edf..baee2d6594 100644
--- a/frontend/src/app/workspace/types/execute-workflow.interface.ts
+++ b/frontend/src/app/workspace/types/execute-workflow.interface.ts
@@ -80,7 +80,6 @@ export enum OperatorState {
export interface OperatorStatistics
extends Readonly<{
- operatorState: OperatorState;
aggregatedInputRowCount: number;
aggregatedInputSize?: number;
inputPortMetrics: Record<string, number>;
@@ -93,9 +92,21 @@ export interface OperatorStatistics
aggregatedIdleTime?: number;
}> {}
+/**
+ * Wire shape of one operator's entry in OperatorStatisticsUpdateEvent. The
+ * engine streams the operator's execution state and its statistics bundled in
+ * one object; WorkflowStatusService splits them into the two separate
+ * sub-concepts (state and statistics).
+ */
+export interface OperatorRuntimeStatus
+ extends OperatorStatistics,
+ Readonly<{
+ operatorState: OperatorState;
+ }> {}
+
export interface OperatorStatsUpdate
extends Readonly<{
- operatorStatistics: Record<string, OperatorStatistics>;
+ operatorStatistics: Record<string, OperatorRuntimeStatus>;
}> {}
export type PaginationMode = { type: "PaginationMode" };