Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,7 @@ import { DashboardWorkflow } from "../../../dashboard/type/dashboard-workflow.in
import { DefaultView } from "../../../dashboard/type/workflow-metadata.interface";
import { SearchFilterParameters, toQueryStrings } from "../../../dashboard/type/search-filter-parameters";
import { NotificationService } from "../notification/notification.service";
import { WorkflowActionService } from "../../../workspace/service/workflow-graph/model/workflow-action.service";
import { last } from "rxjs/operators";

describe("WorkflowPersistService", () => {
Expand All @@ -70,9 +71,15 @@ describe("WorkflowPersistService", () => {
'{"linkID":"link-c94e24a6-2c77-40cf-ba22-1a7ffba64b7d","source":{"operatorID":' +
'"MySQLSource-operator-1ee619b1-8884-4564-a136-29ef77dfcc50","portID":"output-0"},"target":' +
'{"operatorID":"Limit-operator-a11370eb-940a-4f10-8b36-8b413b2396c9","portID":"input-0"}}],"breakpoints":{}}';
// What the page currently holds as the open workflow's metadata (read at response time to keep
// the user's name/description edits). Another workflow by default, so a response is relayed as
// is; the tests about local edits point it at the saved workflow.
let currentMetadata: { wid: number | undefined; name: string; description: string | undefined };
beforeEach(() => {
currentMetadata = { wid: 999, name: "another workflow", description: undefined };
TestBed.configureTestingModule({
imports: [HttpClientTestingModule],
providers: [{ provide: WorkflowActionService, useValue: { getWorkflowMetadata: () => currentMetadata } }],
});
service = TestBed.inject(WorkflowPersistService);
httpTestingController = TestBed.inject(HttpTestingController);
Expand Down Expand Up @@ -260,6 +267,140 @@ describe("WorkflowPersistService", () => {
expect(secondName).toBe("second");
});

describe("a response versus an edit made since the save was sent", () => {
const wf = (name: string) => ({ wid: 9, name, description: "d1", content: validContent }) as unknown as Workflow;

it("relays the response with the page's current name and description, not the ones it was saved with", () => {
currentMetadata = { wid: 9, name: "renamed meanwhile", description: "described meanwhile" };
let result: Workflow | undefined;
service.persistWorkflow(wf("old")).subscribe(w => (result = w));

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "old", description: "d1", lastModifiedTime: 777, content: "{}" });

// Feeding this back as the metadata keeps the rename; the server-owned fields still arrive.
expect(result?.name).toBe("renamed meanwhile");
expect(result?.description).toBe("described meanwhile");
expect(result?.lastModifiedTime).toBe(777);
});

it("keeps the local name for a workflow the save has just created (the page still holds the default id)", () => {
currentMetadata = { wid: 0, name: "named before the first save answered", description: undefined };
let result: Workflow | undefined;
service.persistWorkflow({ ...wf("Untitled workflow"), wid: 0 } as Workflow).subscribe(w => (result = w));

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 42, name: "Untitled workflow", content: "{}" });

expect(result?.wid).toBe(42);
expect(result?.name).toBe("named before the first save answered");
});

it("leaves the response alone when another workflow is open by the time it answers", () => {
currentMetadata = { wid: 10, name: "the other one", description: undefined };
let result: Workflow | undefined;
service.persistWorkflow(wf("old")).subscribe(w => (result = w));

httpTestingController.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`).flush({ wid: 9, name: "old", content: "{}" });

expect(result?.name).toBe("old");
});

it("leaves the response alone when the page was cleared while its save was out", () => {
// clearWorkflow puts the default metadata back, so the page holds the default id like a
// just-created workflow does; the save went out with the real id, which tells them apart.
currentMetadata = { wid: 0, name: "Untitled Workflow", description: undefined };
let result: Workflow | undefined;
service.persistWorkflow(wf("old")).subscribe(w => (result = w));

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "old", description: "d1", content: "{}" });

expect(result?.name).toBe("old");
expect(result?.description).toBe("d1");
});
});

describe("whenSavesDrained", () => {
const wf = (name: string) => ({ wid: 9, name, description: "", content: validContent }) as unknown as Workflow;

it("emits at once when no save is pending", () => {
let emitted = false;
service.whenSavesDrained().subscribe(() => (emitted = true));
expect(emitted).toBe(true);
});

it("emits only once the last queued save has answered, not when the first has", () => {
service.persistWorkflow(wf("first")).subscribe();
service.persistWorkflow(wf("second")).subscribe();
let emitted = false;
service.whenSavesDrained().subscribe(() => (emitted = true));

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "first", content: "{}" });
expect(emitted).toBe(false); // the second is still out

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "second", content: "{}" });
expect(emitted).toBe(true);
});

it("answers a caller asking from its own save's complete callback once the save queued behind has too", () => {
// How the Form View hand-over uses it: its save completes, it asks there, and a rename's save
// queued behind must have answered before it is told the queue is drained.
let emitted = false;
service.persistWorkflow(wf("switch")).subscribe({
complete: () => service.whenSavesDrained().subscribe(() => (emitted = true)),
});
service.persistWorkflow(wf("rename")).subscribe();

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "switch", content: "{}" });
expect(emitted).toBe(false); // asked, and the rename's save is still out

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "rename", content: "{}" });
expect(emitted).toBe(true);
});

it("answers a call once: a later drain does not reach a caller answered already", () => {
// The hand-over's subscription outlives a refused navigation; a drain caused by some later
// save must not run its callback again and route without a click.
let emissions = 0;
service.persistWorkflow(wf("first")).subscribe();
service.whenSavesDrained().subscribe(() => emissions++);
httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "first", content: "{}" });
expect(emissions).toBe(1);

service.persistWorkflow(wf("second")).subscribe();
httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush({ wid: 9, name: "second", content: "{}" });
expect(emissions).toBe(1);
});

it("counts a failed save as done, so a failure does not hold the drain forever", () => {
service.persistWorkflow(wf("first")).subscribe({ error: () => {} });
let emitted = false;
service.whenSavesDrained().subscribe(() => (emitted = true));

httpTestingController
.expectOne(`${API}/${WORKFLOW_PERSIST_URL}`)
.flush("boom", { status: 500, statusText: "Server Error" });

expect(emitted).toBe(true);
});
});

it("persistWorkflow notifies the user when the workflow is broken but still POSTs", () => {
const errorSpy = vi.spyOn(notificationService, "error").mockImplementation(() => {});
const workflow = {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -18,14 +18,15 @@
*/

import { HttpClient, HttpParams } from "@angular/common/http";
import { Injectable } from "@angular/core";
import { EMPTY, Observable, ReplaySubject, Subject, throwError } from "rxjs";
import { catchError, concatMap, filter, map, tap } from "rxjs/operators";
import { Injectable, Injector } from "@angular/core";
import { EMPTY, Observable, of, ReplaySubject, Subject, throwError } from "rxjs";
import { catchError, concatMap, filter, finalize, map, take, tap } from "rxjs/operators";
import { AppSettings } from "../../app-setting";
import { Workflow, WorkflowContent } from "../../type/workflow";
import { DashboardWorkflow } from "../../../dashboard/type/dashboard-workflow.interface";
import { DefaultView } from "../../../dashboard/type/workflow-metadata.interface";
import { WorkflowUtilService } from "../../../workspace/service/workflow-graph/util/workflow-util.service";
import { WorkflowActionService } from "../../../workspace/service/workflow-graph/model/workflow-action.service";
import { NotificationService } from "../notification/notification.service";
import { SearchFilterParameters, toQueryStrings } from "../../../dashboard/type/search-filter-parameters";
import { User } from "../../type/user";
Expand Down Expand Up @@ -68,29 +69,83 @@ export class WorkflowPersistService {
* before it has landed too. Each request snapshots its payload when asked for; it is sent when its
* turn comes, and its outcome is relayed to that caller alone. A failed save fails its own caller
* and does not hold up the next.
*
* A response is relayed with the page's current name and description in place of its own (see
* withLocalEdits): those are the two fields a user edits, and a response answers the save it was
* sent for, which may be older than an edit made since. Callers feed the response back as the
* workflow's metadata; without this, a rename made while a save was out came back undone.
*/
private readonly persistQueue = new Subject<{ send: Observable<Workflow>; result: Subject<Workflow> }>();
private readonly persistQueue = new Subject<{
send: Observable<Workflow>;
result: Subject<Workflow>;
sentWid: number | undefined;
}>();

/** Saves asked for and not yet answered (or failed); see whenSavesDrained. */
private pendingSaves = 0;
private readonly savesDrained = new Subject<void>();

constructor(
private http: HttpClient,
private notificationService: NotificationService
private notificationService: NotificationService,
// Looked up lazily, at response time: the persist service is also used by the dashboard, where
// no workflow is open and constructing the (graph-owning) action service would be a side effect.
private injector: Injector
) {
this.persistQueue
.pipe(
concatMap(({ send, result }) =>
concatMap(({ send, result, sentWid }) =>
send.pipe(
map(updated => this.withLocalEdits(updated, sentWid)),
tap({
next: updated => result.next(updated),
error: (err: unknown) => result.error(err),
complete: () => result.complete(),
}),
catchError(() => EMPTY)
catchError(() => EMPTY),
finalize(() => this.saveDone())
)
)
)
.subscribe();
}

/**
* Emits once every save asked for so far has been answered or has failed; at once when none is
* pending. For a caller that leaves its view on completion (the Form View switch): its own save
* completing is not enough. A save queued behind it (a rename's, a description's) is still sent
* -- this service outlives the view -- but it answers to the component that asked for it, and a
* component the route has destroyed shows no error and feeds back no response.
*/
public whenSavesDrained(): Observable<void> {
return this.pendingSaves === 0 ? of(undefined) : this.savesDrained.pipe(take(1));
}

private saveDone(): void {
this.pendingSaves -= 1;
if (this.pendingSaves === 0) {
this.savesDrained.next();
}
}

/**
* The response with the page's current name and description: a response carries the values the
* save was sent with, and an edit made since would be undone by feeding them back.
*
* Only for a response that is the page's, which is one whose save went out with the id the page
* still holds: the open workflow's, or the default id of a workflow this very save created and
* the page still holds under it. Left alone otherwise: another workflow is open by now, or the
* page was cleared meanwhile (clearWorkflow puts the default id back, but this save went out with
* the real one). Nothing local belongs to those.
*/
private withLocalEdits(response: Workflow, sentWid: number | undefined): Workflow {
const current = this.injector.get(WorkflowActionService).getWorkflowMetadata();
if (current.wid !== sentWid) {
return response;
}
return { ...response, name: current.name, description: current.description };
}

/**
* persists a workflow to backend database and returns its updated information (e.g., new wid).
* The request is queued behind any save still in flight (see persistQueue); the returned
Expand Down Expand Up @@ -122,7 +177,8 @@ export class WorkflowPersistService {
// Replayed, so a caller that subscribes after the queue has already relayed the outcome (a
// save that was quick, or a synchronous test double) still receives it.
const result = new ReplaySubject<Workflow>(1);
this.persistQueue.next({ send, result });
this.pendingSaves += 1;
this.persistQueue.next({ send, result, sentWid: workflow.wid });
return result.asObservable();
}

Expand Down
86 changes: 85 additions & 1 deletion frontend/src/app/workspace/component/menu/menu.component.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -24,9 +24,10 @@ import { HttpClientTestingModule } from "@angular/common/http/testing";
import { RouterTestingModule } from "@angular/router/testing";
import { NzModalService, NzModalModule, NzModalRef } from "ng-zorro-antd/modal";
import { BehaviorSubject, of, Subject, throwError } from "rxjs";
import { take } from "rxjs/operators";
import { WorkflowResultExportService } from "../../service/workflow-result-export/workflow-result-export.service";

import { MenuComponent } from "./menu.component";
import { HANDOVER_DRAIN_TIMEOUT_MS, MenuComponent } from "./menu.component";
import { WorkflowWebsocketService } from "../../service/workflow-websocket/workflow-websocket.service";
import type { ExecutionDurationUpdateEvent } from "../../types/workflow-websocket.interface";
import { OperatorMetadataService } from "../../service/operator-metadata/operator-metadata.service";
Expand Down Expand Up @@ -241,6 +242,89 @@ describe("MenuComponent", () => {
expect(component.isSaving).toBe(false);
});

it("leaves only once every queued save has landed, not just its own", () => {
// A rename made while the switch's save is out saves through the menu itself and is queued
// behind the switch's save; the hand-over waits for the queue to drain, so that save's error
// or response still reaches this component rather than one the route has destroyed.
component.writeAccess = true;
vi.spyOn(component["workflowActionService"], "getWorkflowMetadata").mockReturnValue({ wid: 7 } as any);
vi.spyOn(workflowPersistService, "persistWorkflow").mockReturnValue(of({ wid: 7, name: "saved" } as any));
vi.spyOn(component["workflowActionService"], "setWorkflowMetadata").mockImplementation(() => {});
const drained$ = new Subject<void>();
vi.spyOn(workflowPersistService, "whenSavesDrained").mockReturnValue(drained$.asObservable());
const navigate = vi.spyOn(component as any, "openFormViewPage").mockImplementation(() => {});

component.onClickOpenFormView();

expect(navigate).not.toHaveBeenCalled(); // its own save is done, another is still queued
drained$.next();
expect(navigate).toHaveBeenCalledWith(7);
expect(component.isSaving).toBe(false);
});

it("saves once more when an edit lands while the hand-over waits for the queue to drain", () => {
// Waiting for a queued save is a second window the page stays editable in, after the one the
// switch's own save opened; an edit made in it is stored before leaving, like one made in the first.
const edits = new Subject<unknown>();
vi.spyOn(component["workflowActionService"], "workflowChanged").mockReturnValue(edits.asObservable());
component.ngOnInit();
component.writeAccess = true;
vi.spyOn(component["workflowActionService"], "getWorkflowMetadata").mockReturnValue({ wid: 7 } as any);
vi.spyOn(component["workflowActionService"], "setWorkflowMetadata").mockImplementation(() => {});
const persistSpy = vi
.spyOn(workflowPersistService, "persistWorkflow")
.mockReturnValue(of({ wid: 7, name: "saved" } as any));
const drained$ = new Subject<void>();
// One emission per call, as the service's whenSavesDrained gives.
vi.spyOn(workflowPersistService, "whenSavesDrained").mockImplementation(() => drained$.pipe(take(1)));
const navigate = vi.spyOn(component as any, "openFormViewPage").mockImplementation(() => {});

component.onClickOpenFormView();
expect(persistSpy).toHaveBeenCalledTimes(1);
edits.next(undefined); // an edit while a queued save is still being waited for
drained$.next();

expect(persistSpy).toHaveBeenCalledTimes(2); // saved once more instead of leaving
expect(navigate).not.toHaveBeenCalled();
drained$.next(); // nothing queued behind the second save
expect(navigate).toHaveBeenCalledTimes(1);
expect(navigate).toHaveBeenCalledWith(7);
expect(component.isSaving).toBe(false);
});

it("leaves after a bound when a save queued behind its own never answers", () => {
// The wait is on a request this component did not make; one that never answers must not hold
// the spinner and the button for good. Past the bound the hand-over leaves as it did before
// the wait existed, and an edit landed meanwhile is not saved again: that save would queue
// behind the request that never answers.
vi.useFakeTimers();
try {
const edits = new Subject<unknown>();
vi.spyOn(component["workflowActionService"], "workflowChanged").mockReturnValue(edits.asObservable());
component.ngOnInit();
component.writeAccess = true;
vi.spyOn(component["workflowActionService"], "getWorkflowMetadata").mockReturnValue({ wid: 7 } as any);
vi.spyOn(component["workflowActionService"], "setWorkflowMetadata").mockImplementation(() => {});
const persistSpy = vi
.spyOn(workflowPersistService, "persistWorkflow")
.mockReturnValue(of({ wid: 7, name: "saved" } as any));
vi.spyOn(workflowPersistService, "whenSavesDrained").mockReturnValue(new Subject<void>().asObservable());
const navigate = vi.spyOn(component as any, "openFormViewPage").mockImplementation(() => {});

component.onClickOpenFormView();
edits.next(undefined); // an edit while the queue is waited for
vi.advanceTimersByTime(HANDOVER_DRAIN_TIMEOUT_MS - 1);
expect(navigate).not.toHaveBeenCalled();
vi.advanceTimersByTime(1);

expect(navigate).toHaveBeenCalledWith(7);
expect(persistSpy).toHaveBeenCalledTimes(1); // not saved again behind a request that never answers
expect(component.isSaving).toBe(false);
} finally {
vi.useRealTimers();
}
});

it("ignores a second click while the hand-over is already in progress", () => {
component.writeAccess = true;
vi.spyOn(component["workflowActionService"], "getWorkflowMetadata").mockReturnValue({ wid: 7 } as any);
Expand Down
Loading
Loading