From ea887c569195b18b57c2e6c24b15477c4e4bf704 Mon Sep 17 00:00:00 2001 From: Martin Vu <127065329+VuMartin@users.noreply.github.com> Date: Thu, 13 Aug 2026 21:36:35 +0000 Subject: [PATCH] fix(agent-service): delete links targeting removed input ports when shrinking input ports (#7349) ### What changes were proposed in this PR? This PR fixes a bug in `WorkflowState.updateOperatorInputPorts` where shrinking an operator's input port count left links targeting removed input ports. The method now removes any links whose target port no longer exists by calling `deleteLink`, preventing dangling links from remaining in the workflow state. A regression test was also added to verify that links targeting deleted input ports are removed when the input port count is reduced. ### Any related issues, documentation, discussions? Closes #7169 ### How was this PR tested? Added a regression test covering the case where an operator's input ports are reduced while a link targets one of the removed ports. Verified with: ```bash bun test --test-name-pattern="removes links targeting input ports that are dropped" bun test ``` Before: image After: image ### Was this PR authored or co-authored using generative AI tooling? (backported from commit 12169c2cfe5f92d0c492ae12d8eb2b53a6fec098) Generated-by: ChatGPT (5.5 mini) --- .../src/agent/workflow-state.test.ts | 228 ++++++++++++++++++ agent-service/src/agent/workflow-state.ts | 9 + 2 files changed, 237 insertions(+) diff --git a/agent-service/src/agent/workflow-state.test.ts b/agent-service/src/agent/workflow-state.test.ts index bffc3769eb7..22e191ab69e 100644 --- a/agent-service/src/agent/workflow-state.test.ts +++ b/agent-service/src/agent/workflow-state.test.ts @@ -174,3 +174,231 @@ describe("WorkflowState - getSubDAG", () => { expect(subDag.operators.map(o => o.operatorID)).toEqual(["op2"]); }); }); +<<<<<<< HEAD:agent-service/src/agent/workflow-state.test.ts +======= + +describe("WorkflowState - toLogicalPlan", () => { + test("produces operators, port-indexed links, and an empty reuse list", () => { + const state = new WorkflowState(); + state.addOperator(makeOperator("op1")); + state.addOperator(makeOperator("op2")); + state.addLink(makeLink("l1", "op1", "op2")); + + const plan = state.toLogicalPlan(); + + expect(plan.operators.map(o => o.operatorID)).toEqual(["op1", "op2"]); + // operatorProperties are spread onto each logical operator; type + ports carry through + expect(plan.operators[0].operatorType).toBe("TestOp"); + expect(plan.links).toEqual([ + { fromOpId: "op1", fromPortId: { id: 0, internal: false }, toOpId: "op2", toPortId: { id: 0, internal: false } }, + ]); + expect(plan.opsToReuseResult).toEqual([]); + }); + + test("resolves the link port indices from the operators' port lists", () => { + const state = new WorkflowState(); + state.addOperator( + makeOperator("src", { + outputPorts: [ + { portID: "output-0", displayName: "Output 0" }, + { portID: "output-1", displayName: "Output 1" }, + ], + }) + ); + state.addOperator( + makeOperator("dst", { + inputPorts: [ + { portID: "input-0", displayName: "Input 0" }, + { portID: "input-1", displayName: "Input 1" }, + ], + }) + ); + state.addLink({ + linkID: "l1", + source: { operatorID: "src", portID: "output-1" }, + target: { operatorID: "dst", portID: "input-1" }, + }); + + expect(state.toLogicalPlan().links[0]).toEqual({ + fromOpId: "src", + fromPortId: { id: 1, internal: false }, + toOpId: "dst", + toPortId: { id: 1, internal: false }, + }); + }); +}); + +describe("WorkflowState - traversal", () => { + // chain: op1 -> op2 -> op3 + function seedChain(): WorkflowState { + const state = new WorkflowState(); + state.addOperator(makeOperator("op1")); + state.addOperator(makeOperator("op2")); + state.addOperator(makeOperator("op3")); + state.addLink(makeLink("l1", "op1", "op2")); + state.addLink(makeLink("l2", "op2", "op3")); + return state; + } + + test("getFrontierOperators returns the leaf operators at depth 1", () => { + // op3 is the only operator that is not the source of a link + expect(seedChain().getFrontierOperators(1)).toEqual(["op3"]); + }); + + test("getFrontierOperators expands one hop upstream per depth, topologically ordered", () => { + const state = seedChain(); + expect(state.getFrontierOperators(2)).toEqual(["op2", "op3"]); + expect(state.getFrontierOperators(3)).toEqual(["op1", "op2", "op3"]); + }); + + test("getFrontierOperators returns an empty list for an empty workflow", () => { + expect(new WorkflowState().getFrontierOperators(3)).toEqual([]); + }); + + test("getSubDAG collects the target and every upstream operator and link", () => { + const sub = seedChain().getSubDAG("op3"); + expect(sub.operators.map(o => o.operatorID).sort()).toEqual(["op1", "op2", "op3"]); + expect(sub.links.map(l => l.linkID).sort()).toEqual(["l1", "l2"]); + }); +}); + +describe("WorkflowState - validation state", () => { + const err = (message: string): ValidationError => ({ isValid: false, messages: { general: message } }); + + test("setValidationError stores and getValidationOutput aggregates; clearValidationError removes", () => { + const state = new WorkflowState(); + + state.setValidationError("op1", err("bad")); + expect(state.getValidationOutput().errors).toEqual({ op1: err("bad") }); + + state.clearValidationError("op1"); + expect(state.getValidationOutput().errors).toEqual({}); + }); + + test("setAllValidationErrors replaces the map and recomputes the empty-state flag", () => { + const state = new WorkflowState(); + state.addOperator(makeOperator("op1")); + + state.setAllValidationErrors({ op1: err("x") }); + + const out = state.getValidationOutput(); + expect(out.errors).toEqual({ op1: err("x") }); + expect(out.workflowEmpty).toBe(false); // one enabled operator present + }); + + test("workflowEmpty is true with no operators and when every operator is disabled", () => { + const state = new WorkflowState(); + + state.setAllValidationErrors({}); + expect(state.getValidationOutput().workflowEmpty).toBe(true); + + state.addOperator(makeOperator("op1", { isDisabled: true })); + state.setAllValidationErrors({}); + expect(state.getValidationOutput().workflowEmpty).toBe(true); + }); + + test("getValidationChangedStream emits the aggregated output on every mutation", () => { + const state = new WorkflowState(); + const emitted: Array> = []; + state.getValidationChangedStream().subscribe(v => emitted.push(v)); + + state.setValidationError("op1", err("x")); + state.clearValidationError("op1"); + + expect(emitted).toHaveLength(2); + expect(emitted[0].errors).toEqual({ op1: err("x") }); + expect(emitted[1].errors).toEqual({}); + }); +}); + +describe("WorkflowState - updateOperatorInputPorts", () => { + test("rebuilds the input ports to the requested count, flagging the extras dynamic", () => { + const state = new WorkflowState(); + state.addOperator(makeOperator("op1")); + + expect(state.updateOperatorInputPorts("op1", 3)).toBe(true); + + const ports = state.getOperator("op1")!.inputPorts; + expect(ports.map(p => p.portID)).toEqual(["input-0", "input-1", "input-2"]); + expect(ports[0].isDynamicPort).toBe(false); // port 0 is static + expect(ports[1].isDynamicPort).toBe(true); // extra ports are dynamic + }); + + test("reducing the count drops the extra ports", () => { + const state = new WorkflowState(); + state.addOperator( + makeOperator("op1", { + inputPorts: [ + { portID: "input-0", displayName: "Input 0" }, + { portID: "input-1", displayName: "Input 1" }, + ], + }) + ); + + state.updateOperatorInputPorts("op1", 1); + + expect(state.getOperator("op1")!.inputPorts.map(p => p.portID)).toEqual(["input-0"]); + }); + + test("returns false for a missing operator", () => { + expect(new WorkflowState().updateOperatorInputPorts("missing", 2)).toBe(false); + }); + + test("removes links targeting input ports that are dropped", () => { + const state = new WorkflowState(); + + state.addOperator(makeOperator("src")); + state.addOperator( + makeOperator("op1", { + inputPorts: [ + { portID: "input-0", displayName: "Input 0" }, + { portID: "input-1", displayName: "Input 1" }, + ], + }) + ); + + state.addLink({ + linkID: "l1", + source: { operatorID: "src", portID: "output-0" }, + target: { operatorID: "op1", portID: "input-1" }, + }); + + state.updateOperatorInputPorts("op1", 1); + + expect(state.getAllLinks().map(l => l.linkID)).toEqual([]); + }); +}); + +describe("WorkflowState - workflow content round-trip", () => { + test("setWorkflowContent replaces the state and getWorkflowContent reflects it", () => { + const state = new WorkflowState(); + state.addOperator(makeOperator("stale")); // must be cleared by setWorkflowContent + + const content: WorkflowContent = { + operators: [makeOperator("op1"), makeOperator("op2")], + operatorPositions: { op1: { x: 1, y: 2 } }, + links: [makeLink("l1", "op1", "op2")], + commentBoxes: [], + settings: { dataTransferBatchSize: 123 }, + }; + state.setWorkflowContent(content); + + const out = state.getWorkflowContent(); + expect(out.operators.map(o => o.operatorID)).toEqual(["op1", "op2"]); + expect(out.links.map(l => l.linkID)).toEqual(["l1"]); + expect(out.operatorPositions).toEqual({ op1: { x: 1, y: 2 } }); + expect(out.settings).toEqual({ dataTransferBatchSize: 123 }); + expect(state.getOperator("stale")).toBeUndefined(); + }); + + test("setWorkflowContent falls back to defaults when settings/commentBoxes are absent", () => { + const state = new WorkflowState(); + // The runtime guards against missing optional fields even though the type requires them. + state.setWorkflowContent({ operators: [], operatorPositions: {}, links: [] } as unknown as WorkflowContent); + + const out = state.getWorkflowContent(); + expect(out.settings).toEqual({ dataTransferBatchSize: 400 }); // DEFAULT_WORKFLOW_SETTINGS + expect(out.commentBoxes).toEqual([]); + }); +}); +>>>>>>> 12169c2cf (fix(agent-service): delete links targeting removed input ports when shrinking input ports (#7349)):agent-service/src/agent/workflow-state.spec.ts diff --git a/agent-service/src/agent/workflow-state.ts b/agent-service/src/agent/workflow-state.ts index 04ad2b0e4e8..4e78ee88e58 100644 --- a/agent-service/src/agent/workflow-state.ts +++ b/agent-service/src/agent/workflow-state.ts @@ -189,6 +189,15 @@ export class WorkflowState { inputPorts: newInputPorts, }; this.operators.set(operatorId, updatedOperator); + + const validInputPorts = new Set(newInputPorts.map(port => port.portID)); + + for (const link of this.getAllLinks()) { + if (link.target.operatorID === operatorId && !validInputPorts.has(link.target.portID)) { + this.deleteLink(link.linkID); + } + } + this.operatorPropertyChangeSubject.next({ operator: updatedOperator }); return true; }