This is an automated email from the ASF dual-hosted git repository. github-merge-queue[bot] pushed a commit to branch gh-readonly-queue/main/pr-7911-54a9a597d87eb63304af77c67ab744da37086d03 in repository https://gitbox.apache.org/repos/asf/texera.git
commit 9bb955ccd3082e1a8f43920113f8a82f5266451c Author: Meng Wang <[email protected]> AuthorDate: Tue Aug 25 00:48:18 2026 +0000 test(frontend): cover the agent service's websocket lifecycle and message dispatch (#7911) ### What changes were proposed in this PR? Extends `agent.service.spec.ts` over the WebSocket lifecycle the issue lists — connection setup, the `handleWebSocketMessage` dispatch arms, and teardown — with 15 tests. They live in the existing `websocket lifecycle and events` block and reuse its `FakeWebSocket` double; no new harness. Measured locally with `--coverage`: | `agent.service.ts` | Before | After | | --- | --- | --- | | lines | 301/329 (91.49 %) | **306/329 (93.01 %)** | | branches | 147/186 (79.03 %) | **160/186 (86.02 %)** | | functions | 90/103 | **91/103** | **Connection setup (`startStatePolling`)** - The `wss:`/`ws:` ternary (411): `window.location` is swapped for one test with the `withLocation` helper `virtual-environment.service.spec.ts` already uses (jsdom's own is not writable), and the URL the double was constructed with is asserted in full. The `ws:` side was already covered. - The parse-failure path (424): a payload that is not JSON is logged and does not throw out of the handler; the connection is left intact. - The `onerror` handler (428-429), previously an uncalled function. - The close guard (435), both ways: a stale socket that reports an abnormal close after being swapped out must not clear the newer connection or flip the state. **Message dispatch (`handleWebSocketMessage`)** - `message.workflowContent` (466) — a snapshot with no workflow content leaves the workflow stream and `wsWorkflowActive` alone, so DB polling stays in charge. - `message.step` (478) — a step event with no step changes nothing. - `message.state` (517) — a status event with no state changes nothing. - The error arm (534) — an event with no `error` falls back to `"Agent error occurred"`; the explicit-message side was already covered. - The `default` arm (451/539) — an unrecognised type warns and changes nothing. **Teardown** - `stopStatePolling` (548) both ways, through `deleteAgent`: for a tracked agent the socket is closed and `stopPolling$` completes; for an id that was never tracked it is a no-op. - `activateAgent` (580) — an already-open socket is kept (no second socket) and a socket that is no longer `OPEN` is replaced. `deactivateAgent` always clears the socket, so the inactive-but-still-connected shape this guard defends against is set up directly on the tracking entry, with a comment saying so. - `deactivateAgent` (594-606) — both early returns (unknown agent; already inactive) and the no-socket side of the close branch, reached the ordinary way by letting a normal close event drop the socket first. **One arm is not reachable.** The issue asks for a step "that does not" carry an id (498). `convertApiReActStep` returns ``id: apiStep.id || `${apiStep.messageId}-${apiStep.stepId || 0}` `` — the fallback is a template literal, which always yields a non-empty string, so `convertedStep.id` is never falsy and `if (convertedStep.id)` cannot take its else branch. The test covers the fallback side instead (a step with no id still advances HEAD to `m9-0`) and records why the arm stays half-taken. The remaining uncovered lines are the state accessors and HTTP failure paths, which #7908 covers; the two PRs touch disjoint regions of the file. No production code was changed. ### Any related issues, documentation, discussions? Closes #7888. ### How was this PR tested? `ng test --watch=false --include src/app/workspace/service/agent/agent.service.spec.ts` — 62 passed (47 before, 15 new), repeated 3× for stability. `yarn format:ci` clean. Failure path verified by breaking one assertion in each of the 15 new tests: 16 failed (the 15 plus one pre-existing test that shares the mutated assertion), non-zero exit, then restored to green. Determinism: every socket event is delivered by invoking the double's `onmessage`/`onerror`/`onclose` directly — no real socket, no reconnect timer, and no new `vi.useFakeTimers()`. `withLocation` restores `window.location` in a `finally`, and `vi.restoreAllMocks()` was added to the block's `afterEach` so the `console` spies cannot leak into a later test even if an assertion throws. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (Opus 4.8 [1M context]) --- .../workspace/service/agent/agent.service.spec.ts | 243 +++++++++++++++++++++ 1 file changed, 243 insertions(+) diff --git a/frontend/src/app/workspace/service/agent/agent.service.spec.ts b/frontend/src/app/workspace/service/agent/agent.service.spec.ts index ac98d65c46..62c790ca68 100644 --- a/frontend/src/app/workspace/service/agent/agent.service.spec.ts +++ b/frontend/src/app/workspace/service/agent/agent.service.spec.ts @@ -448,6 +448,7 @@ describe("AgentService", () => { afterEach(() => { vi.unstubAllGlobals(); vi.useRealTimers(); + vi.restoreAllMocks(); }); describe("activateAgent / deactivateAgent", () => { @@ -497,6 +498,135 @@ describe("AgentService", () => { expect(service.isAgentActivelyConnected("agent-1")).toBe(true); expect(service.getActivelyConnectedAgentIds()).toEqual(["agent-1"]); }); + + it("keeps an already-open socket on re-activation and replaces a non-open one", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const tracking = (service as any).agentStateTracking.get("agent-1"); + const ws = FakeWebSocket.latest(); + + // deactivateAgent always drops the socket, so the "inactive but still + // holding a socket" shape that activateAgent's readyState guard defends + // against is set up directly. + ws.readyState = FakeWebSocket.OPEN; + tracking.isActive = false; + expect(service.activateAgent("agent-1")).toBe(true); + expect(FakeWebSocket.instances.length).toBe(1); + expect(tracking.websocket).toBe(ws); + + // The same shape with a socket that is no longer OPEN reconnects instead. + ws.readyState = FakeWebSocket.CLOSED; + tracking.isActive = false; + expect(service.activateAgent("agent-1")).toBe(true); + expect(FakeWebSocket.instances.length).toBe(2); + expect(tracking.websocket).toBe(FakeWebSocket.latest()); + }); + + it("ignores deactivation of an unknown agent and of an already-inactive one", () => { + service.deactivateAgent("nope"); + expect((service as any).agentStateTracking.has("nope")).toBe(false); + + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const ws = FakeWebSocket.latest(); + service.deactivateAgent("agent-1"); + expect(ws.close).toHaveBeenCalledTimes(1); + + const tracking = (service as any).agentStateTracking.get("agent-1"); + const stopPolling = tracking.stopPolling$; + service.deactivateAgent("agent-1"); + + // The second call returns at the isActive guard: nothing is torn down twice. + expect(ws.close).toHaveBeenCalledTimes(1); + expect(tracking.stopPolling$).toBe(stopPolling); + }); + + it("deactivates cleanly when the socket was already dropped by a close event", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const ws = FakeWebSocket.latest(); + const tracking = (service as any).agentStateTracking.get("agent-1"); + + // A normal close clears tracking.websocket but leaves the agent active. + ws.onclose!({ code: 1000 }); + expect(tracking.websocket).toBeUndefined(); + expect(tracking.isActive).toBe(true); + + const stopPolling = tracking.stopPolling$; + service.deactivateAgent("agent-1"); + + expect(ws.close).not.toHaveBeenCalled(); + expect(tracking.isActive).toBe(false); + expect(tracking.stopPolling$).not.toBe(stopPolling); + }); + }); + + describe("connection setup", () => { + /** Swap window.location for the duration of fn; jsdom's own is not writable. */ + const withLocation = <T>(overrides: Partial<Location>, fn: () => T): T => { + const original = window.location; + Object.defineProperty(window, "location", { + configurable: true, + value: { ...original, ...overrides }, + }); + try { + return fn(); + } finally { + Object.defineProperty(window, "location", { configurable: true, value: original }); + } + }; + + it("switches to the wss scheme when the page is served over https", () => { + seedAgent("agent-1"); + + withLocation({ protocol: "https:", host: "texera.example.org" }, () => service.activateAgent("agent-1")); + + expect(FakeWebSocket.latest().url).toBe("wss://texera.example.org/api/agents/agent-1/react"); + }); + + it("logs a payload that is not valid JSON instead of throwing out of the handler", () => { + const errSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const ws = FakeWebSocket.latest(); + + expect(() => ws.onmessage!({ data: "<html>not json</html>" })).not.toThrow(); + + expect(errSpy).toHaveBeenCalledWith("Failed to parse agent WebSocket message:", expect.any(SyntaxError)); + // The connection is left intact. + expect((service as any).agentStateTracking.get("agent-1").websocket).toBe(ws); + }); + + it("logs socket transport errors", () => { + const errSpy = vi.spyOn(console, "error").mockImplementation(() => {}); + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const event = { type: "error" }; + + FakeWebSocket.latest().onerror!(event); + + expect(errSpy).toHaveBeenCalledWith("Agent agent-1 WebSocket error:", event); + }); + + it("ignores a close event from a socket that has already been replaced", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const stale = FakeWebSocket.latest(); + const states: AgentState[] = []; + service.getAgentStateObservable("agent-1").subscribe(s => states.push(s)); + + service.deactivateAgent("agent-1"); + service.activateAgent("agent-1"); + const current = FakeWebSocket.latest(); + expect(current).not.toBe(stale); + + // The stale socket reports an abnormal close after being swapped out: + // the live connection survives and the state does not flip. + stale.onclose!({ code: 1006 }); + + expect((service as any).agentStateTracking.get("agent-1").websocket).toBe(current); + expect(states).toEqual([AgentState.AVAILABLE]); + }); }); describe("sendMessage", () => { @@ -693,6 +823,119 @@ describe("AgentService", () => { FakeWebSocket.latest().onclose!({ code: 1006 }); expect(states[states.length - 1]).toBe(AgentState.UNAVAILABLE); }); + + it("leaves the workflow stream alone when a snapshot carries no workflow content", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const states: AgentState[] = []; + service.getAgentStateObservable("agent-1").subscribe(s => states.push(s)); + const workflows: (Workflow | null)[] = []; + service.getWorkflowObservable("agent-1").subscribe(w => workflows.push(w)); + + emit(FakeWebSocket.latest(), { type: "WsServerSnapshotEvent", state: "GENERATING" }); + + expect(states[states.length - 1]).toBe(AgentState.GENERATING); + expect(workflows).toEqual([null]); + // Workflow polling stays in charge while the socket sends no content. + expect((service as any).agentStateTracking.get("agent-1").wsWorkflowActive).toBe(false); + }); + + it("ignores a step event that carries no step", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + let steps: ReActStep[] = []; + service.getReActStepsObservable("agent-1").subscribe(s => (steps = s)); + + emit(FakeWebSocket.latest(), { type: "WsServerStepEvent" }); + + expect(steps).toEqual([]); + expect(service.getHeadId("agent-1")).toBeNull(); + }); + + it("advances HEAD to the synthesised id when the step carries none", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + + emit(FakeWebSocket.latest(), { + type: "WsServerStepEvent", + step: { messageId: "m9", timestamp: "2026-06-11T00:00:00.000Z" }, + }); + + // convertApiReActStep falls back to `${messageId}-${stepId || 0}`, which is + // always a non-empty string, so a step never fails the HEAD guard. + expect(service.getHeadId("agent-1")).toBe("m9-0"); + }); + + it("ignores a status event that carries no state", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const states: AgentState[] = []; + service.getAgentStateObservable("agent-1").subscribe(s => states.push(s)); + + emit(FakeWebSocket.latest(), { type: "WsServerStatusEvent" }); + + expect(states).toEqual([AgentState.AVAILABLE]); + }); + + it("falls back to a generic notification when the error event carries no message", () => { + vi.spyOn(console, "error").mockImplementation(() => {}); + seedAgent("agent-1"); + service.activateAgent("agent-1"); + + emit(FakeWebSocket.latest(), { type: "WsServerErrorEvent" }); + + expect(notification.error).toHaveBeenCalledWith("Agent error occurred"); + expect((service as any).agents.has("agent-1")).toBe(true); + }); + + it("warns about an unrecognised message type and changes nothing", () => { + const warnSpy = vi.spyOn(console, "warn").mockImplementation(() => {}); + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const states: AgentState[] = []; + service.getAgentStateObservable("agent-1").subscribe(s => states.push(s)); + let steps: ReActStep[] = []; + service.getReActStepsObservable("agent-1").subscribe(s => (steps = s)); + + emit(FakeWebSocket.latest(), { type: "WsServerSomethingElseEvent", state: "GENERATING" }); + + expect(warnSpy).toHaveBeenCalledWith("Unknown agent WebSocket message type:", "WsServerSomethingElseEvent"); + expect(states).toEqual([AgentState.AVAILABLE]); + expect(steps).toEqual([]); + expect(notification.error).not.toHaveBeenCalled(); + }); + }); + + describe("stopStatePolling", () => { + // stopStatePolling is internal; deleteAgent is the public entry point that + // reaches it for both a tracked and an untracked agent. + it("closes the socket and drops the tracking entry for a tracked agent", () => { + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const ws = FakeWebSocket.latest(); + const tracking = (service as any).agentStateTracking.get("agent-1"); + let stopped = false; + tracking.stopPolling$.subscribe({ complete: () => (stopped = true) }); + + let deleted: boolean | undefined; + service.deleteAgent("agent-1").subscribe(d => (deleted = d)); + httpMock.expectOne(r => r.method === "DELETE" && r.url === "/api/agents/agent-1").flush({ deleted: true }); + + expect(deleted).toBe(true); + expect(ws.close).toHaveBeenCalledTimes(1); + expect(tracking.websocket).toBeUndefined(); + expect(stopped).toBe(true); + expect((service as any).agentStateTracking.has("agent-1")).toBe(false); + }); + + it("is a no-op for an agent that was never tracked", () => { + let deleted: boolean | undefined; + service.deleteAgent("ghost").subscribe(d => (deleted = d)); + httpMock.expectOne(r => r.method === "DELETE" && r.url === "/api/agents/ghost").flush({ deleted: true }); + + expect(deleted).toBe(true); + expect((service as any).agentStateTracking.has("ghost")).toBe(false); + }); }); describe("workflow polling", () => {
