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));

Reply via email to