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-8586-5264df2861840bf8b1662ff2e1814b8aa80215db in repository https://gitbox.apache.org/repos/asf/texera.git
commit e7d1676e1662879eb210a21acab8d30779b01024 Author: Meng Wang <[email protected]> AuthorDate: Fri Sep 18 20:53:24 2026 +0000 feat: require a warehouse for every execution while the feature is enabled (#8586) ### What changes were proposed in this PR? With per-user warehouses enabled, an execution carrying no `warehouseId` silently wrote into the shared default warehouse: `resolveLakekeeperWarehouseName` mapped `None` to `None` whenever the flag was on. That made "a run writes into the user's own warehouse" a UI convention rather than a system property — the frontend gates Run on a pick (#7817), but nothing enforced it below, and one path already violated it: agent-driven runs hardcoded `warehouseId = None`, bypassing the picker entirely. **The requirement.** `resolveLakekeeperWarehouseName` refuses an execution that names no warehouse while the feature is enabled, the same way an execution needs a computing unit. Not a privilege change: an explicit `whid` was, and still is, checked against the caller's `uid`. It is resolved before anything destructive happens — both in `initExecutionService`, ahead of the teardown of the execution in flight, and in the sync endpoint, ahead of its own shutdown — so a request that will be refused never disturbs a run already going. **The agent path**, which had no way to carry a pick at all. `SyncExecutionRequest` accepts a `warehouseId` and `SyncExecutionResource` forwards it; the workspace attaches the current pick to each prompt, the agent service applies it before running, and the execute request carries it to the backend. Per prompt rather than at agent creation on purpose: an agent created before the picker had loaded would otherwise carry no warehouse for its whole life, with every run refused and nothing in the agent panel able to correct it. It also means a warehouse chosen after the agent exists takes effect, and an absent pick clears a previous one rather than leaving a stale id to be sent. With the feature off (the default) nothing changes: no pick still means the shared default warehouse, and an explicit pick is still refused loudly (#6930). ### Any related issues, documentation, discussions? Closes #7751. Part of #6870; the last of its Phase 0 items, on top of the dashboard tab (#8005) and the on-canvas picker (#8551). ### How was this PR tested? - Scala: `WorkflowServiceWarehouseSpec` (5) covers the requirement alongside the existing ownership and flag-off cases, and `WorkflowServiceSpec` (11) pins the new ordering — a refused request leaves the previous execution attached. `sbt 'WorkflowExecutionService/testOnly *WorkflowServiceSpec *WorkflowServiceWarehouseSpec'` - `agent-service`: 309 tests (`bun test`, on CI's bun 1.3.3), covering the delegate update and the execute request's body in both the picked and unpicked cases. - Frontend: 87 tests in `agent.service.spec.ts`, pinning the pick attached to each prompt. - Failure paths verified rather than assumed: the requirement, the pre-teardown ordering, the per-prompt carry, the delegate update and the request field were each broken on purpose and the suites confirmed to fail for the expected reason before being restored. - `scalafmtCheck` (both source sets), `tsc --noEmit` and `prettier --check` for `agent-service`, eslint for the frontend files: all clean. ### Was this PR authored or co-authored using generative AI tooling? Generated-by: Claude Code (claude-opus-5, claude-fable-5) --- agent-service/src/agent/texera-agent.spec.ts | 22 ++++++++++++- agent-service/src/agent/texera-agent.ts | 24 +++++++++++++- .../agent/tools/workflow-execution-tools.spec.ts | 26 +++++++++++++++ .../src/agent/tools/workflow-execution-tools.ts | 4 +++ agent-service/src/server.ts | 9 ++++++ agent-service/src/types/agent.ts | 5 +++ agent-service/src/types/ws/client.ts | 7 +++- .../web/resource/SyncExecutionResource.scala | 16 ++++++++-- .../texera/web/service/WorkflowService.scala | 30 +++++++++++++----- .../texera/web/service/WorkflowServiceSpec.scala | 37 ++++++++++++++++++++-- .../web/service/WorkflowServiceWarehouseSpec.scala | 13 ++++++-- .../workspace/service/agent/agent.service.spec.ts | 21 ++++++++++++ .../app/workspace/service/agent/agent.service.ts | 13 ++++++-- 13 files changed, 207 insertions(+), 20 deletions(-) diff --git a/agent-service/src/agent/texera-agent.spec.ts b/agent-service/src/agent/texera-agent.spec.ts index 52d9103866..7c99ffaa3e 100644 --- a/agent-service/src/agent/texera-agent.spec.ts +++ b/agent-service/src/agent/texera-agent.spec.ts @@ -1074,10 +1074,29 @@ describe("delegate mode", () => { } }); + test("setDelegateWarehouse points an existing delegate at the current pick", async () => { + // An agent created before the picker loaded carries no warehouse; every run + // would be refused, and nothing in the agent panel could fix it (#7751). + const agent = makeAgentWith(textModel("x")); + agent.setDelegateWarehouse(42); + expect((agent as any).delegateConfig).toBeUndefined(); + + agent.setDelegateConfig({ userToken: "tok", workflowId: 7 }); + expect((agent as any).buildExecutionConfig().warehouseId).toBeUndefined(); + + agent.setDelegateWarehouse(42); + expect((agent as any).buildExecutionConfig().warehouseId).toBe(42); + + // An absent pick is a pick: clearing it keeps a stale id from riding the + // next run and being refused while the feature is off. + agent.setDelegateWarehouse(undefined); + expect((agent as any).buildExecutionConfig().warehouseId).toBeUndefined(); + }); + test("buildExecutionConfig projects the delegate config and live settings", async () => { const agent = makeAgentWith(textModel("x")); expect((agent as any).buildExecutionConfig()).toBeUndefined(); - (agent as any).delegateConfig = { userToken: "tok", workflowId: 5, computingUnitId: 2 }; + (agent as any).delegateConfig = { userToken: "tok", workflowId: 5, computingUnitId: 2, warehouseId: 42 }; agent.updateSettings({ executionTimeoutMs: 7000, maxOperatorResultCharLimit: 11, @@ -1087,6 +1106,7 @@ describe("delegate mode", () => { userToken: "tok", workflowId: 5, computingUnitId: 2, + warehouseId: 42, maxOperatorResultCharLimit: 11, maxOperatorResultCellCharLimit: 13, executionTimeoutMs: 7000, diff --git a/agent-service/src/agent/texera-agent.ts b/agent-service/src/agent/texera-agent.ts index 9a640aaab6..34b5f32f98 100644 --- a/agent-service/src/agent/texera-agent.ts +++ b/agent-service/src/agent/texera-agent.ts @@ -111,6 +111,7 @@ export class TexeraAgent { workflowId: number; workflowName?: string; computingUnitId?: number; + warehouseId?: number; }; private stepCallback: ReActStepCallback | null = null; @@ -185,6 +186,7 @@ export class TexeraAgent { userToken: this.delegateConfig.userToken, workflowId: this.delegateConfig.workflowId, computingUnitId: this.delegateConfig.computingUnitId, + warehouseId: this.delegateConfig.warehouseId, maxOperatorResultCharLimit: this.settings.maxOperatorResultCharLimit, maxOperatorResultCellCharLimit: this.settings.maxOperatorResultCellCharLimit, executionTimeoutMs: this.settings.executionTimeoutMs, @@ -425,6 +427,7 @@ export class TexeraAgent { workflowId: number; workflowName?: string; computingUnitId?: number; + warehouseId?: number; }): void { this.delegateConfig = config; @@ -433,8 +436,27 @@ export class TexeraAgent { this.setupWorkflowChangeHandlers(); } + /** + * Point the delegate at the warehouse the workspace has selected now. The rest + * of the config is fixed at creation; this one travels per prompt because the + * user can pick (or first load) a warehouse after the agent exists (#7751). + */ + setDelegateWarehouse(warehouseId: number | undefined): void { + if (!this.delegateConfig || this.delegateConfig.warehouseId === warehouseId) { + return; + } + this.delegateConfig = { ...this.delegateConfig, warehouseId }; + } + getDelegateConfig(): - | { userToken: string; userInfo?: UserInfo; workflowId: number; workflowName?: string; computingUnitId?: number } + | { + userToken: string; + userInfo?: UserInfo; + workflowId: number; + workflowName?: string; + computingUnitId?: number; + warehouseId?: number; + } | undefined { return this.delegateConfig; } diff --git a/agent-service/src/agent/tools/workflow-execution-tools.spec.ts b/agent-service/src/agent/tools/workflow-execution-tools.spec.ts index 396ef15bf8..361337af44 100644 --- a/agent-service/src/agent/tools/workflow-execution-tools.spec.ts +++ b/agent-service/src/agent/tools/workflow-execution-tools.spec.ts @@ -370,6 +370,32 @@ describe("executeOperatorAndFormat — request construction", () => { ); }); + test("the execute request carries the picked warehouse, and omits it when none is picked", async () => { + // The last hop of the chain: the id has travelled request -> agent -> + // delegate config -> here, and the backend refuses a run without it while + // the feature is enabled (#7751). + const state = new WorkflowState(); + state.addOperator(makeOperator("dst")); + resolveFetch(fetchSpy, { + success: true, + state: "Completed", + operators: { dst: { state: "Completed", inputTuples: 0, outputTuples: 1, resultMode: "table", result: [] } }, + }); + + await executeOperatorAndFormat(state, cfg({ workflowId: 7, computingUnitId: 3, warehouseId: 42 }), "dst"); + expect(requestBody(fetchSpy).warehouseId).toBe(42); + + fetchSpy.mockClear(); + resolveFetch(fetchSpy, { + success: true, + state: "Completed", + operators: { dst: { state: "Completed", inputTuples: 0, outputTuples: 1, resultMode: "table", result: [] } }, + }); + + await executeOperatorAndFormat(state, cfg({ workflowId: 7, computingUnitId: 3 }), "dst"); + expect("warehouseId" in requestBody(fetchSpy)).toBe(false); + }); + test("sends the whole workflow, not an upstream slice, when the operator id is empty", async () => { // The plan builder only takes the sub-DAG path for a *truthy* target id, so an empty // operator id falls through to the "every operator" branch. `sink` and `orphan` are the diff --git a/agent-service/src/agent/tools/workflow-execution-tools.ts b/agent-service/src/agent/tools/workflow-execution-tools.ts index 13c5a24158..efe9f2f57c 100644 --- a/agent-service/src/agent/tools/workflow-execution-tools.ts +++ b/agent-service/src/agent/tools/workflow-execution-tools.ts @@ -37,6 +37,9 @@ export interface ExecutionConfig { userToken: string; workflowId: number; computingUnitId?: number; + // The warehouse the user picked in the UI; forwarded alongside computingUnitId so + // the run writes into their own warehouse rather than shared storage (#7751). + warehouseId?: number; maxOperatorResultCharLimit?: number; maxOperatorResultCellCharLimit?: number; executionTimeoutMs?: number; @@ -272,6 +275,7 @@ async function executeWorkflowHttp( maxOperatorResultCharLimit: config.maxOperatorResultCharLimit ?? DEFAULT_AGENT_SETTINGS.maxOperatorResultCharLimit, maxOperatorResultCellCharLimit: config.maxOperatorResultCellCharLimit ?? DEFAULT_AGENT_SETTINGS.maxOperatorResultCellCharLimit, + ...(config.warehouseId !== undefined ? { warehouseId: config.warehouseId } : {}), }; log.debug( diff --git a/agent-service/src/server.ts b/agent-service/src/server.ts index 030d27b95b..42247399e8 100644 --- a/agent-service/src/server.ts +++ b/agent-service/src/server.ts @@ -130,6 +130,7 @@ function getAgentInfo(agentId: string, agent: TexeraAgent): AgentInfo { workflowId: delegateConfig.workflowId, workflowName: delegateConfig.workflowName, computingUnitId: delegateConfig.computingUnitId, + warehouseId: delegateConfig.warehouseId, } : undefined, settings: settingsApi, @@ -492,6 +493,14 @@ export function buildApp() { wsLog.info({ agentId, preview: msg.content.substring(0, 50) }, "received command"); + // The prompt carries the workspace's current warehouse pick, so a run + // uses what the user has selected now rather than whatever was + // selected when the agent was created. An absent field IS the + // current selection — none — so it clears a previous pick rather + // than leaving a stale id to be sent (and refused while the feature + // is off) forever (#7751). + agent.setDelegateWarehouse(typeof msg.warehouseId === "number" ? msg.warehouseId : undefined); + agent.setStepCallback((step: ReActStep) => { broadcastToAgentClients(agentId, new WsServerStepEvent(step)); }); diff --git a/agent-service/src/types/agent.ts b/agent-service/src/types/agent.ts index 694b51785f..89bf75aa40 100644 --- a/agent-service/src/types/agent.ts +++ b/agent-service/src/types/agent.ts @@ -121,6 +121,11 @@ export interface AgentDelegateConfig { workflowId?: number; workflowName?: string; computingUnitId?: number; + // The warehouse the delegating user has picked in the workspace. Unlike the + // rest of this config it is not fixed at creation: each prompt carries the + // current pick, so a run always writes into what the user has selected now + // (#7751). + warehouseId?: number; } export interface AgentSettingsApi { diff --git a/agent-service/src/types/ws/client.ts b/agent-service/src/types/ws/client.ts index 2983ed1d6c..ba03ac829c 100644 --- a/agent-service/src/types/ws/client.ts +++ b/agent-service/src/types/ws/client.ts @@ -27,7 +27,12 @@ export class WsClientPromptCommand { readonly type = "WsClientPromptCommand"; constructor( readonly content: string, - readonly messageSource?: "chat" | "feedback" + readonly messageSource?: "chat" | "feedback", + // The warehouse picked in the workspace right now. Sent per prompt rather + // than fixed at agent creation: an agent created before the picker loaded + // would otherwise carry no warehouse for its whole life and every run would + // be refused, with nothing in the agent panel able to fix it (#7751). + readonly warehouseId?: number ) {} } diff --git a/amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala b/amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala index cd528a6cf5..9073f37ade 100644 --- a/amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala +++ b/amber/src/main/scala/org/apache/texera/web/resource/SyncExecutionResource.scala @@ -73,7 +73,12 @@ case class SyncExecutionRequest( targetOperatorIds: List[String], timeoutSeconds: Int, maxOperatorResultCharLimit: Int, - maxOperatorResultCellCharLimit: Int + maxOperatorResultCellCharLimit: Int, + // The user_warehouse this run writes into. Carried from the caller the same way + // computingUnitId is: the agent forwards what the user picked in the UI. Required + // while per-user warehouses are enabled; absent keeps the shared default when the + // feature is off (#7751). + warehouseId: Option[Int] = None ) case class ConsoleMessageInfo( @@ -152,6 +157,13 @@ class SyncExecutionResource extends LazyLogging { computingUnitId ) + // Same rule as initExecutionService: check the pick before anything + // destructive, or an agent request that is going to be refused takes the + // execution in flight down with it. Resolved again inside the init call — + // one indexed single-row read, against a request that already writes + // several rows. + WorkflowService.resolveLakekeeperWarehouseName(request.warehouseId, user.getUser.getUid) + shutdownPreviousExecution(workflowService) // "Execute To" semantics: when a single target is given, run only its upstream sub-DAG. @@ -169,7 +181,7 @@ class SyncExecutionResource extends LazyLogging { ), emailNotificationEnabled = false, computingUnitId = computingUnitId, - warehouseId = None + warehouseId = request.warehouseId ) workflowService.initExecutionService( diff --git a/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala b/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala index 02394bee97..8ec650ba71 100644 --- a/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala +++ b/amber/src/main/scala/org/apache/texera/web/service/WorkflowService.scala @@ -72,9 +72,15 @@ object WorkflowService { /** * Maps an execution's chosen warehouse (its user_warehouse row id) to its Lakekeeper - * warehouse name, checking that the requesting user owns it. `None` (no explicit pick) keeps the - * shared default warehouse. With warehouses disabled, an explicit pick is refused - * loudly rather than silently routed into the shared warehouse (#6930). + * warehouse name, checking that the requesting user owns it. + * + * With warehouses enabled a pick is **required**, the same way a computing unit is: + * falling back to the shared default would make "a run writes into the user's own + * warehouse" a UI convention rather than a system property, and would silently route + * any caller that forgot to pick into shared storage (#7751). + * + * With warehouses disabled, an explicit pick is refused loudly rather than silently + * routed into the shared warehouse, and no pick keeps the shared default (#6930). */ def resolveLakekeeperWarehouseName( warehouseId: Option[Int], @@ -89,6 +95,11 @@ object WorkflowService { ) return None } + if (warehouseId.isEmpty) { + throw new IllegalArgumentException( + "a warehouse must be selected for this execution" + ) + } warehouseId.map(id => { val row = SqlServer .getInstance() @@ -219,21 +230,24 @@ class WorkflowService( sessionUri: URI ): Unit = { - if (executionService.hasValue) { - executionService.getValue.unsubscribeAll() - } - val (uidOpt, userEmailOpt) = userOpt.map(user => (user.getUid, user.getEmail)).unzip + // Validate before touching the execution already in flight: a request that is + // going to be refused must not take the running one's subscriptions with it. // uid is NOT NULL in the DB; fail early here rather than letting the insert fail downstream. val uid = uidOpt.getOrElse( throw new IllegalArgumentException( "Cannot start execution: a user id (uid) is required but none was provided." ) ) + val warehouseName = WorkflowService.resolveLakekeeperWarehouseName(req.warehouseId, uid) + + if (executionService.hasValue) { + executionService.getValue.unsubscribeAll() + } val workflowContext: WorkflowContext = createWorkflowContext() - workflowContext.warehouse = WorkflowService.resolveLakekeeperWarehouseName(req.warehouseId, uid) + workflowContext.warehouse = warehouseName var coordinatorConf = CoordinatorConfig.default // clean up results from previous run diff --git a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceSpec.scala index 210d821b87..9e80cb2989 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceSpec.scala @@ -595,7 +595,37 @@ class WorkflowServiceSpec events shouldBe empty } - it should "refuse to start an execution with no user id, after detaching the previous one" in { + it should "refuse an invalid warehouse pick before detaching the previous execution" in { + // The pick is resolved before the teardown: an agent (or a stale client) + // asking for a warehouse the deployment will reject must not take the + // execution already running down with it (#7751). With the feature off, + // any explicit pick is refused (#6930), which is the cheapest way to make + // resolution fail here. + val service = new TestWorkflowService(9L) + val previous = newExecution() + service.executionService.onNext(previous) + val events = collectExecutionEvents(previous) + val request = WorkflowExecuteRequest( + executionName = "test", + engineVersion = "test", + logicalPlan = LogicalPlanPojo(List.empty, List.empty, List.empty, List.empty), + replayFromExecution = None, + workflowSettings = WorkflowSettings(), + emailNotificationEnabled = false, + computingUnitId = 1, + warehouseId = Some(1) + ) + + val error = intercept[IllegalArgumentException] { + service.initExecutionService(request, Some(executingUser), new URI("vfs:///session")) + } + error.getMessage should include("warehouse") + + previous.executionStateStore.metadataStore.updateState(_.withState(RUNNING)) + events should not be empty + } + + it should "refuse to start an execution with no user id, leaving the previous one attached" in { val service = new TestWorkflowService(6L) val previous = newExecution() service.executionService.onNext(previous) @@ -617,9 +647,10 @@ class WorkflowServiceSpec } error.getMessage should include("user id") - // The previous execution is detached first, whatever happens next. + // A request that is going to be refused must not take the running execution + // with it: its subscriptions are still live. previous.executionStateStore.metadataStore.updateState(_.withState(RUNNING)) - events shouldBe empty + events should not be empty } it should "clear the previous run's storage registry, on its computing unit only, before starting a new execution" in { diff --git a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala index 07c7d45176..788190e9f3 100644 --- a/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala +++ b/amber/src/test/scala/org/apache/texera/web/service/WorkflowServiceWarehouseSpec.scala @@ -72,8 +72,17 @@ class WorkflowServiceWarehouseSpec override protected def afterAll(): Unit = closeConnectionPool() - "resolveLakekeeperWarehouseName" should "keep the shared default warehouse when nothing is picked" in { - WorkflowService.resolveLakekeeperWarehouseName(None, ownerUid, enabled = true) shouldBe None + "resolveLakekeeperWarehouseName" should "require a pick while warehouses are enabled" in { + // A run must name the warehouse it writes into, the same way it names a computing + // unit; falling back to the shared default would route a caller that forgot to pick + // into shared storage (#7751). + val error = intercept[IllegalArgumentException] { + WorkflowService.resolveLakekeeperWarehouseName(None, ownerUid, enabled = true) + } + error.getMessage should include("warehouse") + } + + it should "keep the shared default warehouse when nothing is picked and the feature is off" in { WorkflowService.resolveLakekeeperWarehouseName(None, ownerUid, enabled = false) shouldBe None } 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 62c790ca68..d5e38a9742 100644 --- a/frontend/src/app/workspace/service/agent/agent.service.spec.ts +++ b/frontend/src/app/workspace/service/agent/agent.service.spec.ts @@ -25,6 +25,7 @@ import { AgentState, ReActStep } from "./agent-types"; import { NotificationService } from "../../../common/service/notification/notification.service"; import { WorkflowPersistService } from "../../../common/service/workflow-persist/workflow-persist.service"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; +import { WarehouseService } from "../../../common/service/warehouse/warehouse.service"; import { DashboardWorkflowComputingUnit } from "../../../common/type/workflow-computing-unit"; import { Workflow } from "../../../common/type/workflow"; import { commonTestProviders } from "../../../common/testing/test-utils"; @@ -63,6 +64,7 @@ describe("AgentService", () => { let service: AgentService; let httpMock: HttpTestingController; let selectedUnit: DashboardWorkflowComputingUnit | null; + let selectedWarehouseId: number | undefined; let notification: Record<"error" | "success" | "info" | "warning", ReturnType<typeof vi.fn>>; let workflowPersist: { retrieveWorkflow: ReturnType<typeof vi.fn> }; @@ -102,6 +104,7 @@ describe("AgentService", () => { beforeEach(() => { selectedUnit = null; + selectedWarehouseId = undefined; notification = { error: vi.fn(), success: vi.fn(), info: vi.fn(), warning: vi.fn() }; workflowPersist = { retrieveWorkflow: vi.fn().mockReturnValue(of(stubWorkflow)) }; TestBed.configureTestingModule({ @@ -114,6 +117,10 @@ describe("AgentService", () => { provide: ComputingUnitStatusService, useValue: { getSelectedComputingUnitValue: () => selectedUnit }, }, + { + provide: WarehouseService, + useValue: { getSelectedWarehouseIdValue: () => selectedWarehouseId }, + }, ...commonTestProviders, ], }); @@ -641,6 +648,20 @@ describe("AgentService", () => { expect(notification.error).toHaveBeenCalledWith("WebSocket connection not available"); }); + it("carries the warehouse picked right now, not the one fixed at agent creation", () => { + // An agent created before the picker loaded would otherwise carry no + // warehouse for its whole life and every run would be refused (#7751). + seedAgent("agent-1"); + service.activateAgent("agent-1"); + const ws = FakeWebSocket.latest(); + ws.readyState = FakeWebSocket.OPEN; + selectedWarehouseId = 42; + + service.sendMessage("agent-1", "run it"); + + expect(JSON.parse(ws.send.mock.calls[0][0]).warehouseId).toBe(42); + }); + it("sends a WsClientPromptCommand carrying the message source over an open socket", () => { seedAgent("agent-1"); service.activateAgent("agent-1"); diff --git a/frontend/src/app/workspace/service/agent/agent.service.ts b/frontend/src/app/workspace/service/agent/agent.service.ts index 5e7c254f22..0464ec8ba8 100644 --- a/frontend/src/app/workspace/service/agent/agent.service.ts +++ b/frontend/src/app/workspace/service/agent/agent.service.ts @@ -40,6 +40,7 @@ import { AppSettings } from "../../../common/app-setting"; import { AgentState, ReActStep, ModelMessage } from "./agent-types"; import { Workflow, WorkflowContent } from "../../../common/type/workflow"; import { ComputingUnitStatusService } from "../../../common/service/computing-unit/computing-unit-status/computing-unit-status.service"; +import { WarehouseService } from "../../../common/service/warehouse/warehouse.service"; /** * Agent settings for API (serializable format). @@ -226,7 +227,8 @@ export class AgentService { private notificationService: NotificationService, private workflowPersistService: WorkflowPersistService, private ngZone: NgZone, - private computingUnitStatusService: ComputingUnitStatusService + private computingUnitStatusService: ComputingUnitStatusService, + private warehouseService: WarehouseService ) { // Sync local cache with backend on service initialization // This handles cases where the backend was restarted @@ -869,11 +871,18 @@ export class AgentService { return; } - const wsMessage = { + const wsMessage: { type: string; content: string; messageSource: string; warehouseId?: number } = { type: "WsClientPromptCommand", content: message, messageSource, }; + // Sent per message, not fixed at agent creation: an agent created before the + // warehouse picker loaded would otherwise never carry one, and every run it + // attempted would be refused (#7751). + const selectedWarehouseId = this.warehouseService.getSelectedWarehouseIdValue(); + if (selectedWarehouseId !== undefined) { + wsMessage.warehouseId = selectedWarehouseId; + } try { tracking.websocket.send(JSON.stringify(wsMessage));
