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>
This commit is contained in:
Hamza Rashid
2026-04-20 13:04:27 -04:00
committed by GitHub
parent b4f0870980
commit ab70449427
9 changed files with 105 additions and 59 deletions

View File

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

View File

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

View File

@ -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,
*,

View File

@ -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):

View File

@ -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

View File

@ -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),
},
)

View File

@ -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"],
}
]

View File

@ -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

View File

@ -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"}]})
# ---------------------------------------------------------------------------