From ab7044942701a75bb80f6bf77ea8bf89be2ace92 Mon Sep 17 00:00:00 2001 From: Hamza Rashid <74062092+HzaRashid@users.noreply.github.com> Date: Mon, 20 Apr 2026 13:04:27 -0400 Subject: [PATCH] feat(deployments): filter WXO agent list to drafts when load_from_provider=true (#12745) * feat(deployments): filter WXO list to drafts and expose environments as a list When `load_from_provider=true` on `GET /deployments`, restrict watsonx Orchestrate results to draft agents only and surface their environment metadata as `environments: list[str]` on both the list and status responses. - Add `BaseDeploymentMapper.resolve_load_from_provider_deployment_list_params` as an extension point; override in the WXO mapper to always inject `{"environment": "draft"}` when listing from the provider. - In the WXO adapter's `list`, pop `environment` from `provider_params` before building the WXO query (it is not a native WXO query param) and apply it as a client-side membership filter. - Skip forwarding `deployment_type` from the endpoint on the `load_from_provider` path; WXO only exposes agent deployments today and other list logic does not rely on it. - Replace the single `environment` string with `environments: list[str]` in list and status `provider_data`. This is a breaking change to the provider-backed shape, accepted because the feature is still behind a flag. - Introduce `get_agent_environments` with fail-fast access (no silent fallbacks/normalization) so WXO contract breaks surface immediately, and drop the now-unused `derive_agent_environment` categorizer. Tests updated to cover the new filter, the list shape, the status shape, and the fail-fast semantics of the helper. * [autofix.ci] apply automated fixes --------- Co-authored-by: autofix-ci[bot] <114827586+autofix-ci[bot]@users.noreply.github.com> --- .../base/langflow/api/v1/deployments.py | 8 ++-- .../api/v1/mappers/deployments/base.py | 7 +++ .../deployments/watsonx_orchestrate/mapper.py | 4 ++ .../watsonx_orchestrate/payloads.py | 21 ++++++-- .../watsonx_orchestrate/core/status.py | 39 +++++++-------- .../deployment/watsonx_orchestrate/service.py | 16 +++++-- .../api/v1/test_deployment_mapper_watsonx.py | 16 +++++-- .../api/v1/test_deployment_route_handlers.py | 5 ++ .../deployment/test_watsonx_orchestrate.py | 48 +++++++++++-------- 9 files changed, 105 insertions(+), 59 deletions(-) diff --git a/src/backend/base/langflow/api/v1/deployments.py b/src/backend/base/langflow/api/v1/deployments.py index 23be02d17f..cc53922635 100644 --- a/src/backend/base/langflow/api/v1/deployments.py +++ b/src/backend/base/langflow/api/v1/deployments.py @@ -644,9 +644,7 @@ async def list_deployments( if flow_ids: resolved = await flow_version_ids_for_flows(session, flow_ids=flow_ids, user_id=current_user.id) if not resolved: - return DeploymentListResponse( - deployments=[], page=params.page, size=params.size, total=0, deployment_type=deployment_type - ) + return DeploymentListResponse(deployments=[], page=params.page, size=params.size, total=0) effective_flow_version_ids = resolved provider_account = await get_owned_provider_account_or_404( @@ -655,11 +653,13 @@ async def list_deployments( deployment_adapter = resolve_deployment_adapter(provider_account.provider_key) deployment_mapper = get_deployment_mapper(provider_account.provider_key) if load_from_provider: + provider_list_params = deployment_mapper.resolve_load_from_provider_deployment_list_params() + adapter_params = DeploymentListParams(provider_params=provider_list_params) if provider_list_params else None with handle_adapter_errors(mapper=deployment_mapper), deployment_provider_scope(provider_id): provider_view = await deployment_adapter.list( user_id=current_user.id, db=session, - params=None if deployment_type is None else DeploymentListParams(deployment_types=[deployment_type]), + params=adapter_params, ) return deployment_mapper.shape_deployment_list_result(provider_view) diff --git a/src/backend/base/langflow/api/v1/mappers/deployments/base.py b/src/backend/base/langflow/api/v1/mappers/deployments/base.py index b68041808f..ffda59f596 100644 --- a/src/backend/base/langflow/api/v1/mappers/deployments/base.py +++ b/src/backend/base/langflow/api/v1/mappers/deployments/base.py @@ -230,6 +230,13 @@ class BaseDeploymentMapper: ) -> dict[str, Any] | None: return self._validate_slot(self.api_payloads.deployment_list_params, raw) + def resolve_load_from_provider_deployment_list_params(self) -> dict[str, Any] | None: + """Return provider_params for provider-backed deployment listing. + + Default behavior applies no provider-specific filters. + """ + return None + async def resolve_config_list_params(self, raw: dict[str, Any] | None, db: AsyncSession) -> dict[str, Any] | None: return self._validate_slot(self.api_payloads.config_list_params, raw) diff --git a/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/mapper.py b/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/mapper.py index a677ef0d8b..1851289b4d 100644 --- a/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/mapper.py +++ b/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/mapper.py @@ -316,6 +316,10 @@ class WatsonxOrchestrateDeploymentMapper(BaseDeploymentMapper): ) return parsed.model_dump(mode="json", exclude_none=True) + def resolve_load_from_provider_deployment_list_params(self) -> dict[str, Any] | None: + """Force provider-backed list mode to draft agents only.""" + return {"environment": "draft"} + def resolve_credentials( self, *, diff --git a/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/payloads.py b/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/payloads.py index 712a2f22e3..29faf99d1a 100644 --- a/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/payloads.py +++ b/src/backend/base/langflow/api/v1/mappers/deployments/watsonx_orchestrate/payloads.py @@ -479,7 +479,7 @@ class WatsonxApiProviderDeploymentListItem(BaseModel): created_at: datetime | None = None updated_at: datetime | None = None tool_ids: list[str] = Field(default_factory=list) - environment: str | None = None + environments: list[str] = Field(default_factory=list) @field_validator("tool_ids", mode="before") @classmethod @@ -488,11 +488,22 @@ class WatsonxApiProviderDeploymentListItem(BaseModel): return [] return [normalized for tool_id in value if (normalized := str(tool_id).strip())] - @field_validator("environment", mode="before") + @field_validator("environments", mode="before") @classmethod - def normalize_environment(cls, value: Any) -> str | None: - normalized = str(value or "").strip() - return normalized or None + def normalize_environments(cls, value: Any) -> list[str]: + if value is None: + return [] + if not isinstance(value, list): + return [] + seen: set[str] = set() + normalized_names: list[str] = [] + for item in value: + name = str(item or "").strip().lower() + if not name or name in seen: + continue + seen.add(name) + normalized_names.append(name) + return normalized_names class WatsonxApiDeploymentListProviderData(BaseModel): diff --git a/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/status.py b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/status.py index 98ffc06536..f7d2b137e1 100644 --- a/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/status.py +++ b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/status.py @@ -4,7 +4,6 @@ from __future__ import annotations from typing import Any -from ibm_watsonx_orchestrate_core.types.connections import ConnectionEnvironment from lfx.services.adapters.deployment.schema import DeploymentGetResult, DeploymentType, ItemResult @@ -46,27 +45,21 @@ def get_deployment_detail_metadata( return DeploymentGetResult(**result) -def derive_agent_environment(agent: dict[str, Any]) -> str: - environments = agent.get("environments", []) - if not isinstance(environments, list) or not environments: - return "unknown" +def get_agent_environments(agent: dict[str, Any]) -> list[str]: + """Return the de-duplicated list of environment names for an agent. - has_draft = False - has_live = False - for env in environments: - if not isinstance(env, dict): + Names are surfaced as-is from the provider so the API response stays + resilient if watsonx Orchestrate adds new environments in the future. + Missing keys or unexpected types indicate a contract break on watsonx + Orchestrate's side and are allowed to raise. + """ + raw_environments = agent["environments"] + seen: set[str] = set() + names: list[str] = [] + for env in raw_environments: + env_name = env["name"] + if env_name in seen: continue - env_name = str(env.get("name", "")).strip().lower() - if env_name == ConnectionEnvironment.DRAFT.value: - has_draft = True - continue - if env_name: - has_live = True - - if has_draft and has_live: - return "both" - if has_live: - return "live" - if has_draft: - return "draft" - return "unknown" + seen.add(env_name) + names.append(env_name) + return names diff --git a/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/service.py b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/service.py index 1e185ff3e3..2cb68118ea 100644 --- a/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/service.py +++ b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/service.py @@ -86,7 +86,7 @@ from langflow.services.adapters.deployment.watsonx_orchestrate.core.retry import rollback_created_resources, ) from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import ( - derive_agent_environment, + get_agent_environments, get_deployment_detail_metadata, get_deployment_metadata, ) @@ -354,9 +354,15 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): raise InvalidDeploymentTypeError(message=msg) query_params: dict[str, Any] = {} + environment_filter: str | None = None if params and params.provider_params: - query_params = params.provider_params + provider_params = dict(params.provider_params) + environment_raw = provider_params.pop("environment", None) + if environment_raw is not None: + normalized_environment = str(environment_raw).strip().lower() + environment_filter = normalized_environment or None + query_params = provider_params if params and params.deployment_ids and "ids" not in query_params: query_params["ids"] = [str(_id) for _id in params.deployment_ids] @@ -371,13 +377,15 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): client_manager.get_agents_raw, params=query_params or None, ) + if environment_filter is not None: + raw_agents = [agent for agent in raw_agents if environment_filter in get_agent_environments(agent)] deployments = [ get_deployment_metadata( data=agent, deployment_type=DeploymentType.AGENT, provider_data={ "tool_ids": extract_agent_tool_ids(agent), - "environment": derive_agent_environment(agent), + "environments": get_agent_environments(agent), }, ) for agent in raw_agents @@ -619,7 +627,7 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): id=agent_id, provider_data={ "status": "connected", - "environment": derive_agent_environment(agent), + "environments": get_agent_environments(agent), }, ) diff --git a/src/backend/tests/unit/api/v1/test_deployment_mapper_watsonx.py b/src/backend/tests/unit/api/v1/test_deployment_mapper_watsonx.py index 532dcac0d5..965bfd2e14 100644 --- a/src/backend/tests/unit/api/v1/test_deployment_mapper_watsonx.py +++ b/src/backend/tests/unit/api/v1/test_deployment_mapper_watsonx.py @@ -86,6 +86,11 @@ def test_watsonx_mapper_is_registered() -> None: assert mapper.api_payloads.snapshot_list_result is not None +def test_watsonx_mapper_load_from_provider_params_force_draft_filter() -> None: + mapper = WatsonxOrchestrateDeploymentMapper() + assert mapper.resolve_load_from_provider_deployment_list_params() == {"environment": "draft"} + + @pytest.mark.asyncio @pytest.mark.parametrize( "provider_data", @@ -180,7 +185,10 @@ def test_watsonx_mapper_provider_list_entry_flattens_provider_data_and_uses_id() description="desc", created_at=now, updated_at=now, - provider_data={"tool_ids": ["tool-1", " ", "tool-2"], "environment": "draft"}, + provider_data={ + "tool_ids": ["tool-1", " ", "tool-2"], + "environments": ["draft", "draft", "live"], + }, ) shaped = mapper._shape_provider_deployment_list_entry(item) @@ -192,7 +200,7 @@ def test_watsonx_mapper_provider_list_entry_flattens_provider_data_and_uses_id() assert shaped["created_at"] == now.isoformat().replace("+00:00", "Z") assert shaped["updated_at"] == now.isoformat().replace("+00:00", "Z") assert shaped["tool_ids"] == ["tool-1", "tool-2"] - assert shaped["environment"] == "draft" + assert shaped["environments"] == ["draft", "live"] assert "provider_data" not in shaped assert "resource_key" not in shaped @@ -209,7 +217,7 @@ def test_watsonx_mapper_shapes_deployment_list_result_with_flattened_entries() - description="desc", created_at=now, updated_at=now, - provider_data={"tool_ids": ["tool-1"], "environment": "live"}, + provider_data={"tool_ids": ["tool-1"], "environments": ["live"]}, ) ] ) @@ -229,7 +237,7 @@ def test_watsonx_mapper_shapes_deployment_list_result_with_flattened_entries() - "created_at": now.isoformat().replace("+00:00", "Z"), "updated_at": now.isoformat().replace("+00:00", "Z"), "tool_ids": ["tool-1"], - "environment": "live", + "environments": ["live"], } ] diff --git a/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py b/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py index 71543f3a7a..373a7a77d5 100644 --- a/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py +++ b/src/backend/tests/unit/api/v1/test_deployment_route_handlers.py @@ -34,6 +34,7 @@ from lfx.services.adapters.deployment.schema import ( ConfigListItem, ConfigListResult, DeploymentCreateResult, + DeploymentListParams, DeploymentListResult, DeploymentUpdateResult, ItemResult, @@ -566,6 +567,7 @@ class TestListDeploymentsLoadFromProvider: mapper = MagicMock() expected = MagicMock() mapper.shape_deployment_list_result.return_value = expected + mapper.resolve_load_from_provider_deployment_list_params.return_value = {"environment": "draft"} mock_get_mapper.return_value = mapper result = await list_deployments( @@ -581,6 +583,9 @@ class TestListDeploymentsLoadFromProvider: assert result is expected mock_list_synced.assert_not_awaited() adapter.list.assert_awaited_once() + list_call_kwargs = adapter.list.await_args.kwargs + assert isinstance(list_call_kwargs["params"], DeploymentListParams) + assert list_call_kwargs["params"].provider_params == {"environment": "draft"} mapper.shape_deployment_list_result.assert_called_once() @pytest.mark.asyncio diff --git a/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py b/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py index f6ef0893b4..a2a0782989 100644 --- a/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py +++ b/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py @@ -1966,9 +1966,9 @@ async def test_list_deployments_filters_with_provider_draft_filters(monkeypatch) fake_agent = FakeAgentClient( {"id": "dep-1", "tools": []}, listed_agents=[ - {"id": "dep-1", "name": "deployment-1", "tools": []}, - {"id": "dep-2", "name": "deployment-2", "tools": []}, - {"id": "dep-3", "name": "deployment-3", "tools": []}, + {"id": "dep-1", "name": "deployment-1", "tools": [], "environments": [{"name": "draft"}]}, + {"id": "dep-2", "name": "deployment-2", "tools": [], "environments": [{"name": "prod"}]}, + {"id": "dep-3", "name": "deployment-3", "tools": [], "environments": [{"name": "draft"}]}, ], ) fake_clients = _with_wxo_wrappers( @@ -1990,11 +1990,11 @@ async def test_list_deployments_filters_with_provider_draft_filters(monkeypatch) db=object(), params=DeploymentListParams( deployment_types=[DeploymentType.AGENT], - provider_params={"ids": ["dep-2"], "names": ["deployment-3"]}, + provider_params={"ids": ["dep-2"], "names": ["deployment-3"], "environment": "draft"}, ), ) - assert sorted(item.id for item in result.deployments) == ["dep-2", "dep-3"] + assert sorted(item.id for item in result.deployments) == ["dep-3"] @pytest.mark.anyio @@ -4812,6 +4812,7 @@ async def test_get_status_connected(monkeypatch): result = await service.get_status(user_id="user-1", deployment_id="dep-1", db=object()) assert result.id == "dep-1" assert result.provider_data["status"] == "connected" + assert result.provider_data["environments"] == ["draft"] @pytest.mark.anyio @@ -4893,7 +4894,7 @@ async def test_list_deployments_without_params(monkeypatch): fake_agent = FakeAgentClient( {"id": "dep-1", "tools": []}, listed_agents=[ - {"id": "dep-1", "name": "agent-1", "tools": []}, + {"id": "dep-1", "name": "agent-1", "tools": [], "environments": [{"name": "draft"}]}, ], ) fake_clients = _with_wxo_wrappers( @@ -5832,29 +5833,38 @@ def test_create_agent_run_result_omits_thread_id_when_absent(): # --------------------------------------------------------------------------- -def test_derive_agent_environment_draft(): - from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import derive_agent_environment +def test_get_agent_environments_dedupes_preserving_order(): + from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import get_agent_environments - assert derive_agent_environment({"environments": [{"name": "draft"}]}) == "draft" + agent = { + "environments": [ + {"name": "draft"}, + {"name": "draft"}, + {"name": "live"}, + {"name": "future-env"}, + ] + } + assert get_agent_environments(agent) == ["draft", "live", "future-env"] -def test_derive_agent_environment_live(): - from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import derive_agent_environment +def test_get_agent_environments_returns_empty_list_when_provider_returns_empty(): + from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import get_agent_environments - assert derive_agent_environment({"environments": [{"name": "production"}]}) == "live" + assert get_agent_environments({"environments": []}) == [] -def test_derive_agent_environment_both(): - from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import derive_agent_environment +def test_get_agent_environments_raises_when_environments_key_missing(): + from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import get_agent_environments - assert derive_agent_environment({"environments": [{"name": "draft"}, {"name": "prod"}]}) == "both" + with pytest.raises(KeyError): + get_agent_environments({}) -def test_derive_agent_environment_empty(): - from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import derive_agent_environment +def test_get_agent_environments_raises_when_env_entry_missing_name(): + from langflow.services.adapters.deployment.watsonx_orchestrate.core.status import get_agent_environments - assert derive_agent_environment({}) == "unknown" - assert derive_agent_environment({"environments": []}) == "unknown" + with pytest.raises(KeyError): + get_agent_environments({"environments": [{"not_name": "draft"}]}) # ---------------------------------------------------------------------------