From 6d90dc614f59a1d967cb39516c0df55916b294eb Mon Sep 17 00:00:00 2001 From: himavarshagoutham Date: Wed, 1 Apr 2026 15:46:09 -0400 Subject: [PATCH 1/3] feat(deployments): verify wxO credentials against instance API - Probe GET /v1/orchestrate/models after IAM token to validate URL+key - Add POST /deployments/providers/verify-credentials for connection tests - Extend unit tests for models probe and WxO client stub --- .../langflow/api/v1/schemas/deployments.py | 7 ---- .../watsonx_orchestrate/core/models.py | 13 ++++++ .../deployment/watsonx_orchestrate/service.py | 41 +++++++++++++++++-- .../deployment/test_watsonx_orchestrate.py | 38 +++++++++++++++++ 4 files changed, 88 insertions(+), 11 deletions(-) create mode 100644 src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/models.py diff --git a/src/backend/base/langflow/api/v1/schemas/deployments.py b/src/backend/base/langflow/api/v1/schemas/deployments.py index b13eea1bf4..56e6be5dc3 100644 --- a/src/backend/base/langflow/api/v1/schemas/deployments.py +++ b/src/backend/base/langflow/api/v1/schemas/deployments.py @@ -39,9 +39,6 @@ from pydantic import AfterValidator, BaseModel, Field, ValidationInfo, model_val from langflow.services.database.models.deployment_provider_account.schemas import ( DeploymentProviderKey, ) -from langflow.services.database.models.deployment_provider_account.utils import ( - validate_provider_url, -) # --------------------------------------------------------------------------- # Shared validation helpers @@ -82,10 +79,6 @@ NonEmptyStr = Annotated[str, AfterValidator(_strip_nonempty)] """String type that strips whitespace and rejects empty/whitespace-only values.""" -ValidatedUrl = Annotated[str, AfterValidator(validate_provider_url)] -"""URL type that enforces HTTPS and normalizes.""" - - def _validate_flow_version_ids(values: list[UUID] | None) -> list[UUID] | None: """AfterValidator for optional flow_version_ids query parameter.""" if values is None: diff --git a/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/models.py b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/models.py new file mode 100644 index 0000000000..c8da87f811 --- /dev/null +++ b/src/backend/base/langflow/services/adapters/deployment/watsonx_orchestrate/core/models.py @@ -0,0 +1,13 @@ +"""Model-catalog retrieval helpers for the wxO deployment adapter.""" + +from __future__ import annotations + +from typing import TYPE_CHECKING, Any + +if TYPE_CHECKING: + from langflow.services.adapters.deployment.watsonx_orchestrate.types import WxOClient + + +def fetch_models_adapter(clients: WxOClient, params: dict[str, Any] | None = None) -> Any: + """Fetch raw provider models through the adapter client seam.""" + return clients.get_models_raw(params=params) 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 fbb223bda8..a85200e953 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 @@ -77,6 +77,9 @@ from langflow.services.adapters.deployment.watsonx_orchestrate.core.execution im create_agent_run, get_agent_run, ) +from langflow.services.adapters.deployment.watsonx_orchestrate.core.models import ( + fetch_models_adapter, +) from langflow.services.adapters.deployment.watsonx_orchestrate.core.retry import ( retry_create, rollback_created_resources, @@ -106,6 +109,7 @@ from langflow.services.adapters.deployment.watsonx_orchestrate.payloads import ( WatsonxDeploymentUpdatePayload, WatsonxDeploymentUpdateResultData, ) +from langflow.services.adapters.deployment.watsonx_orchestrate.types import WxOClient from langflow.services.adapters.deployment.watsonx_orchestrate.utils import ( dedupe_list, extract_agent_tool_ids, @@ -124,8 +128,6 @@ if TYPE_CHECKING: from lfx.services.settings.service import SettingsService from sqlalchemy.ext.asyncio import AsyncSession - from langflow.services.adapters.deployment.watsonx_orchestrate.types import WxOClient - class WatsonxOrchestrateDeploymentService(BaseDeploymentService): """Deployment adapter for Watsonx Orchestrate.""" @@ -296,7 +298,7 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): """List provider-available LLM model names.""" client_manager = await self._get_provider_clients(user_id=user_id, db=db) try: - raw_models = await asyncio.to_thread(client_manager.get_models_raw) + raw_models = await asyncio.to_thread(fetch_models_adapter, client_manager) parsed_models: WatsonxDeploymentLlmListResultData = self._parse_provider_payload( slot=self.payload_schemas.deployment_llm_list_result, slot_name="deployment_llm_list_result", @@ -830,7 +832,13 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): user_id: IdLike, # noqa: ARG002 payload: VerifyCredentials, ) -> VerifyCredentialsResult: - """Verify WXO credentials by obtaining a token from the provider.""" + """Verify WXO credentials for the target instance. + + Obtains an IAM/MCSP token, then calls the wxO models listing API for the + configured instance URL. Token-only checks are insufficient because a + valid API key may authenticate while still lacking access to the tenant + represented by the instance URL. + """ verify_slot = self.payload_schemas.verify_credentials if verify_slot is None: msg = "Required slot 'verify_credentials' is not configured." @@ -883,6 +891,31 @@ class WatsonxOrchestrateDeploymentService(BaseDeploymentService): cause=exc, ) from exc + def _probe_instance_models() -> None: + wxo_client = WxOClient(instance_url=payload.base_url, authenticator=authenticator) + fetch_models_adapter(wxo_client) + + try: + await asyncio.to_thread(_probe_instance_models) + except ClientAPIException as exc: + status_code = exc.response.status_code if exc.response is not None else None + logger.error( # noqa: TRY400 + "Credential verification failed: wxO instance probe rejected request (status=%s)", + status_code, + ) + raise_deployment_error_from_status( + status_code=status_code, + detail="Credential verification failed.", + message_prefix="Credential verification", + cause=None, + ) + except Exception as exc: + raise DeploymentError( + message="Credential verification failed unexpectedly.", + error_code="deployment_error", + cause=exc, + ) from exc + return VerifyCredentialsResult() async def update_snapshot( 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 0dbab60d32..94b32ed889 100644 --- a/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py +++ b/src/backend/tests/unit/services/deployment/test_watsonx_orchestrate.py @@ -6864,6 +6864,7 @@ async def test_verify_credentials_success(monkeypatch): "get_authenticator", lambda **_kwargs: FakeAuthenticator(), ) + monkeypatch.setattr(service_module, "fetch_models_adapter", lambda *_args, **_kwargs: {}) svc = WatsonxOrchestrateDeploymentService(settings_service=DummySettingsService()) payload = VerifyCredentials( @@ -6903,6 +6904,43 @@ async def test_verify_credentials_invalid_key_raises(monkeypatch): await svc.verify_credentials(user_id="u1", payload=payload) +@pytest.mark.anyio +async def test_verify_credentials_instance_probe_forbidden(monkeypatch): + """403 from wxO models probe maps to AuthorizationError (wrong instance for key).""" + from ibm_watsonx_orchestrate_clients.tools.tool_client import ClientAPIException + from lfx.services.adapters.deployment.exceptions import AuthorizationError + from lfx.services.adapters.deployment.schema import VerifyCredentials + from requests import Response + + class FakeTokenManager: + def get_token(self): + return "fake-token" + + class FakeAuthenticator: + token_manager = FakeTokenManager() + + monkeypatch.setattr( + service_module, + "get_authenticator", + lambda **_kwargs: FakeAuthenticator(), + ) + + def _fail_fetch_models(*_args, **_kwargs): + response = Response() + response.status_code = 403 + raise ClientAPIException(response=response) + + monkeypatch.setattr(service_module, "fetch_models_adapter", _fail_fetch_models) + + svc = WatsonxOrchestrateDeploymentService(settings_service=DummySettingsService()) + payload = VerifyCredentials( + base_url="https://api.us-south.wxo.cloud.ibm.com", + provider_data={"api_key": "valid-key"}, # pragma: allowlist secret + ) + with pytest.raises(AuthorizationError, match="Credential verification"): + await svc.verify_credentials(user_id="u1", payload=payload) + + @pytest.mark.anyio async def test_verify_credentials_malformed_key_from_authenticator_constructor_raises(): """verify_credentials raises InvalidContentError when authenticator creation fails validation.""" From 5ea740ca56e8f3852df87fac9f665f8ec1880506 Mon Sep 17 00:00:00 2001 From: Hamza Rashid Date: Thu, 2 Apr 2026 21:50:26 +0000 Subject: [PATCH 2/3] fix(api): rename verify endpoint and redact 422 request input Rename the deployment provider credential verification route to /deployments/providers/verify and sanitize RequestValidationError responses so raw request payloads are not echoed back in 422 errors. Update route/tests accordingly, including regression coverage for input redaction. --- src/backend/base/langflow/main.py | 15 ++++++++++++++- src/backend/tests/unit/api/v1/test_projects.py | 18 ++++++++++++++++++ 2 files changed, 32 insertions(+), 1 deletion(-) diff --git a/src/backend/base/langflow/main.py b/src/backend/base/langflow/main.py index fd1d3db7dc..c13e4f4d22 100644 --- a/src/backend/base/langflow/main.py +++ b/src/backend/base/langflow/main.py @@ -8,13 +8,14 @@ import warnings from contextlib import asynccontextmanager, suppress from http import HTTPStatus from pathlib import Path -from typing import TYPE_CHECKING, cast +from typing import TYPE_CHECKING, Any, cast from urllib.parse import urlencode import anyio import httpx import sqlalchemy from fastapi import FastAPI, HTTPException, Request, Response, status +from fastapi.exceptions import RequestValidationError from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import FileResponse, JSONResponse from fastapi.staticfiles import StaticFiles @@ -78,6 +79,11 @@ async def log_exception_to_telemetry(exc: Exception, context: str) -> None: await logger.awarning(f"Failed to log {context} exception to telemetry") +def _sanitize_validation_errors(errors: list[dict[str, Any]]) -> list[dict[str, Any]]: + """Strip request payload echoes from validation errors to avoid leaking submitted data.""" + return [{key: value for key, value in error.items() if key != "input"} for error in errors] + + class RequestCancelledMiddleware(BaseHTTPMiddleware): def __init__(self, app) -> None: super().__init__(app) @@ -560,6 +566,13 @@ def create_app(): # Discover and register additional routers from plugins (langflow.plugins entry-point) load_plugin_routes(app) + @app.exception_handler(RequestValidationError) + async def request_validation_exception_handler(_request: Request, exc: RequestValidationError): + return JSONResponse( + status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, + content={"detail": _sanitize_validation_errors(exc.errors())}, + ) + @app.exception_handler(Exception) async def exception_handler(_request: Request, exc: Exception): if isinstance(exc, HTTPException): diff --git a/src/backend/tests/unit/api/v1/test_projects.py b/src/backend/tests/unit/api/v1/test_projects.py index c48c06baa0..8ddaf9678f 100644 --- a/src/backend/tests/unit/api/v1/test_projects.py +++ b/src/backend/tests/unit/api/v1/test_projects.py @@ -134,6 +134,24 @@ async def test_create_project_validation_error(client: AsyncClient, logged_in_he assert response.status_code == status.HTTP_422_UNPROCESSABLE_ENTITY +async def test_create_project_validation_error_does_not_echo_request_input(client: AsyncClient, logged_in_headers): + """Regression test: request-validation 422 responses must not echo submitted payloads.""" + sensitive_marker = "secret-api-key-should-not-echo" + response = await client.post( + "api/v1/projects/", + # Send a JSON string (instead of object) to trigger RequestValidationError consistently. + content=json.dumps(sensitive_marker), + headers={**logged_in_headers, "Content-Type": "application/json"}, + ) + + assert response.status_code == status.HTTP_422_UNPROCESSABLE_ENTITY + detail = response.json()["detail"] + assert isinstance(detail, list) + assert detail + assert all("input" not in error for error in detail) + assert sensitive_marker not in response.text + + async def test_delete_project_then_404(client: AsyncClient, logged_in_headers, basic_case): create_resp = await client.post("api/v1/projects/", json=basic_case, headers=logged_in_headers) proj_id = create_resp.json()["id"] From af325e9205374990a6efdd285509cbbca6d26221 Mon Sep 17 00:00:00 2001 From: himavarshagoutham Date: Tue, 14 Apr 2026 09:58:10 -0400 Subject: [PATCH 3/3] revert(api): remove verification endpoint follow-up changes Keep this PR focused on wxO adapter tenant validation via model-list fetch and drop API-layer validation endpoint and request-redaction changes. --- .../langflow/api/v1/schemas/deployments.py | 7 +++++++ src/backend/base/langflow/main.py | 15 +-------------- src/backend/tests/unit/api/v1/test_projects.py | 18 ------------------ 3 files changed, 8 insertions(+), 32 deletions(-) diff --git a/src/backend/base/langflow/api/v1/schemas/deployments.py b/src/backend/base/langflow/api/v1/schemas/deployments.py index 56e6be5dc3..b13eea1bf4 100644 --- a/src/backend/base/langflow/api/v1/schemas/deployments.py +++ b/src/backend/base/langflow/api/v1/schemas/deployments.py @@ -39,6 +39,9 @@ from pydantic import AfterValidator, BaseModel, Field, ValidationInfo, model_val from langflow.services.database.models.deployment_provider_account.schemas import ( DeploymentProviderKey, ) +from langflow.services.database.models.deployment_provider_account.utils import ( + validate_provider_url, +) # --------------------------------------------------------------------------- # Shared validation helpers @@ -79,6 +82,10 @@ NonEmptyStr = Annotated[str, AfterValidator(_strip_nonempty)] """String type that strips whitespace and rejects empty/whitespace-only values.""" +ValidatedUrl = Annotated[str, AfterValidator(validate_provider_url)] +"""URL type that enforces HTTPS and normalizes.""" + + def _validate_flow_version_ids(values: list[UUID] | None) -> list[UUID] | None: """AfterValidator for optional flow_version_ids query parameter.""" if values is None: diff --git a/src/backend/base/langflow/main.py b/src/backend/base/langflow/main.py index c13e4f4d22..fd1d3db7dc 100644 --- a/src/backend/base/langflow/main.py +++ b/src/backend/base/langflow/main.py @@ -8,14 +8,13 @@ import warnings from contextlib import asynccontextmanager, suppress from http import HTTPStatus from pathlib import Path -from typing import TYPE_CHECKING, Any, cast +from typing import TYPE_CHECKING, cast from urllib.parse import urlencode import anyio import httpx import sqlalchemy from fastapi import FastAPI, HTTPException, Request, Response, status -from fastapi.exceptions import RequestValidationError from fastapi.middleware.cors import CORSMiddleware from fastapi.responses import FileResponse, JSONResponse from fastapi.staticfiles import StaticFiles @@ -79,11 +78,6 @@ async def log_exception_to_telemetry(exc: Exception, context: str) -> None: await logger.awarning(f"Failed to log {context} exception to telemetry") -def _sanitize_validation_errors(errors: list[dict[str, Any]]) -> list[dict[str, Any]]: - """Strip request payload echoes from validation errors to avoid leaking submitted data.""" - return [{key: value for key, value in error.items() if key != "input"} for error in errors] - - class RequestCancelledMiddleware(BaseHTTPMiddleware): def __init__(self, app) -> None: super().__init__(app) @@ -566,13 +560,6 @@ def create_app(): # Discover and register additional routers from plugins (langflow.plugins entry-point) load_plugin_routes(app) - @app.exception_handler(RequestValidationError) - async def request_validation_exception_handler(_request: Request, exc: RequestValidationError): - return JSONResponse( - status_code=status.HTTP_422_UNPROCESSABLE_ENTITY, - content={"detail": _sanitize_validation_errors(exc.errors())}, - ) - @app.exception_handler(Exception) async def exception_handler(_request: Request, exc: Exception): if isinstance(exc, HTTPException): diff --git a/src/backend/tests/unit/api/v1/test_projects.py b/src/backend/tests/unit/api/v1/test_projects.py index 8ddaf9678f..c48c06baa0 100644 --- a/src/backend/tests/unit/api/v1/test_projects.py +++ b/src/backend/tests/unit/api/v1/test_projects.py @@ -134,24 +134,6 @@ async def test_create_project_validation_error(client: AsyncClient, logged_in_he assert response.status_code == status.HTTP_422_UNPROCESSABLE_ENTITY -async def test_create_project_validation_error_does_not_echo_request_input(client: AsyncClient, logged_in_headers): - """Regression test: request-validation 422 responses must not echo submitted payloads.""" - sensitive_marker = "secret-api-key-should-not-echo" - response = await client.post( - "api/v1/projects/", - # Send a JSON string (instead of object) to trigger RequestValidationError consistently. - content=json.dumps(sensitive_marker), - headers={**logged_in_headers, "Content-Type": "application/json"}, - ) - - assert response.status_code == status.HTTP_422_UNPROCESSABLE_ENTITY - detail = response.json()["detail"] - assert isinstance(detail, list) - assert detail - assert all("input" not in error for error in detail) - assert sensitive_marker not in response.text - - async def test_delete_project_then_404(client: AsyncClient, logged_in_headers, basic_case): create_resp = await client.post("api/v1/projects/", json=basic_case, headers=logged_in_headers) proj_id = create_resp.json()["id"]