Repository navigation
Add a Databricks toolset for Unity Gateway MCP Services - #74198
firasbouzazi wants to merge 7 commits into
Conversation
|
Congratulations on your first Pull Request and welcome to the Apache Airflow community! If you have any issues or are unsure about any anything please check our Contributors' Guide
|
|
Hello @shahar1 can you assaign a reviewer please. |
|
cc @zozo123 and @moomindani for review |
zozo123
left a comment
There was a problem hiding this comment.
Thanks @firasbouzazi, and welcome! This follows #74193 closely: Databricks auth and addressing stay in the databricks provider, common.ai is optional, the token never leaves the workspace host, and the stub-gateway tests run the real MCP client. The URL shape and the EXECUTE + USE CATALOG/USE SCHEMA requirement match the Unity Gateway docs, and test_unity_mcp.py passes locally (31 tests). I also ran the toolset read-only against a real workspace with a PAT. get_tools on system.ai.dbsql works, a tool-side error comes back as ModelRetry as intended, and the token never shows up in logs.
A few things before this can go in:
P1
- Retry vs. fatal comes from state shared across calls.
_DatabricksTokenAuthkeeps oneerror_statusper toolset, and pydantic-ai runs tool calls concurrently. I ran two calls in parallel, one getting a 503 and the other a 200. In 4 of 20 runs the failed call came back asModelRetry, so the model gets invited to re-run a call that may already have applied. That's the one thing #74193 says must not happen, andsystem.ai.dbsql'sexecute_sqlcan write. The same design also turns a400that carries a JSON-RPCInvalid paramsbody into a fatalDatabricksUnityMCPError. Databricks' own tools return bad arguments as 200isError(I checked), but registered external servers may well send a 400. Suggestion: classify from the per-call exception. If the chain has theMcpErrorthe MCP client synthesizes for an HTTP status ("Server returned an error response", "Not Found", "Session terminated"), raiseDatabricksUnityMCPError, always fortools/call. If it's a real JSON-RPC error from the server, let it stay aModelRetry. Keep the recorded status only to enrich the message (401/403/429, Retry-After). Please add a concurrency test and a 400-with-JSON-RPC-body test. - CI.
- Static checks:
check_doc_filesonly accepts how-to guides underdocs/operators|sensors|transfer. Drop thehow-to-guidefrom the newDatabricks Unity Gatewayintegration inprovider.yamland regenerateget_provider_info.py. - Spellcheck: AutoAPI renders
example_databricks_unity_mcp, so add lowercasemcptodocs/spelling_wordlist.txt.
- Static checks:
P2
3. Coupling to common.ai internals. This overrides MCPToolset._get_server, sets _server/_tool_prefix, and passes the Databricks conn id as mcp_conn_id. common.ai's stability guarantee only covers MCPToolset's public parameters, and this would be the first toolset subclass outside common.ai, so any internal refactor there breaks Databricks. I'd rather keep it KISS:
- either add a small public per-request auth hook to common.ai's MCPHook/MCPToolset (an
httpx2.Auth, or atoken_providerthat's called per request) in its own PR and compose with it here, - or subclass
AirflowToolsetand own thePydanticAIMCPToolsetdirectly.
Placement in the databricks provider looks right to me either way.
4. A missing service is a 403, not a 404. On a live workspace, a nonexistent service (main.default.does_not_exist) returns HTTP 403 with JSON-RPC -32007 "Not authorized to invoke MCP service.". So DatabricksUnityMCPServiceNotFoundError never fires for that case, and the errors table in the docs is off. I'd drop the class, or keep it only for a bare 404, and word the 403 message as "doesn't exist, or the identity lacks EXECUTE / USE CATALOG / USE SCHEMA".
5. Proxies. The proxies extra on the Databricks connection is ignored by the MCP transport. Please pass an httpx_client_factory with the proxy, or document it as unsupported.
6. Dev dependency. pyproject.toml dev group: use apache-airflow-providers-common-ai[mcp]. Without it, uv sync in providers/databricks doesn't install fastmcp and the whole test module silently skips.
7. Live run. #74193 asked for live validation to be recorded separately from the mocks. I covered PAT. Could you note in the PR whether you ran the system test, and whether service principal OAuth or Azure worked? Note that the hook requests scope=all-apis, so an SP secret scoped only to ai-gateway will fail.
8. GenAI disclosure. Small process thing: the branch is claude/.... If AI tooling helped, please tick the GenAI box and add the Generated-by: line per the PR guidelines.
P3
_endpoint_urlhonours the connectionschema, sohttpwould send the bearer token in cleartext. Consider requiring https except on loopback.- Docs:
- Please add a source for "account users already hold these privileges on the built-in
system.aiservices", and for Entra/AAD tokens being accepted by Unity Gateway. - Mention that a service principal's OAuth secret must allow
all-apis. The hook requestsall-apis, so a secret scoped only toai-gatewayfails.
- Please add a source for "account users already hold these privileges on the built-in
import httpx2relies on a transitive dependency. It's fine, just flagging it.
Happy to re-review once (1)–(3) are in. Thanks for picking this up!
| self.error_status = self.retry_after = self.token_error = None | ||
| return error | ||
|
|
||
| def _record(self, response: httpx2.Response) -> None: |
There was a problem hiding this comment.
This status is shared across concurrent calls. Another call's 2xx can clear it before the failed call reads it, and then an ambiguous 5xx on tools/call surfaces as ModelRetry (I reproduced 4 of 20 with two parallel calls). Could the fatal-vs-retry decision come from the call's own exception chain (the client's synthesized McpError), with this kept only as a hint for the message?
There was a problem hiding this comment.
there's no shared error state any more: the fatal-vs-retry decision comes from the call's own exception chain. A client-synthesized McpError on tools/call raises DatabricksUnityMCPError and is never retried. The HTTP status is still recorded, but only for JSON-RPC requests, and it's only used for the message and exception type when the operation ran alone; with parallel calls it's left out. test_concurrent_calls_are_each_judged_by_their_own_failure holds two calls at a barrier so they really overlap (503 + 200, 503 + 429, 400-with-JSON-RPC + 503).
| ambiguous = ( | ||
| " The tool call may or may not have run, so it was not retried." if during_tool_call else "" | ||
| ) | ||
| if status is not None: |
There was a problem hiding this comment.
A 400 that carries a JSON-RPC error body (e.g. -32602 Invalid params) ends up here as a fatal error, so the model can't correct its arguments. The MCP client already surfaces those bodies as a real McpError. Let them pass through as ModelRetry.
There was a problem hiding this comment.
A real JSON-RPC error from the server now passes through as ModelRetry at any HTTP status, so the model can correct its arguments. Only the client's stand-in errors (no JSON-RPC answer) are fatal. test_json_rpc_error_from_the_server_is_left_to_the_model covers 200, 400 and 500 with a -32602 body.
| "USE CATALOG and USE SCHEMA on its catalog and schema, and its credentials must be valid.", | ||
| http_status_code=status, | ||
| ) | ||
| if status == 404: |
There was a problem hiding this comment.
On a live workspace a nonexistent service returns 403 with JSON-RPC -32007 Not authorized to invoke MCP service., not a 404, so this branch doesn't fire for a missing service. Maybe fold it into the 403 message ('doesn't exist, or lacks EXECUTE...') and fix the docs table.
| raise ValueError(f"Connection {self._databricks_conn_id!r} has no workspace host.") | ||
| return hook._endpoint_url(f"{GATEWAY_MCP_SERVICES_PATH}/{self._service_name}") | ||
|
|
||
| def _get_server(self) -> Any: |
There was a problem hiding this comment.
This overrides private MCPToolset internals (_get_server, _server, _tool_prefix) from another provider. Could we either add a public per-request auth hook to common.ai, or subclass AirflowToolset directly?
There was a problem hiding this comment.
Went with subclassing AirflowToolset directly in 6402af7: the toolset owns its PydanticAIMCPToolset and no longer touches _get_server, _server or _tool_prefix of MCPToolset. This also showed the extra's floor was wrong. In common.ai 0.10.0, MCPToolset.call_tool didn't go through execute_tool, so the error handling would never have run. The extra now requires >=0.11.0. A public per-request auth hook in common.ai could still be a follow-up if you think it's worth it.
| # Fail here, with a clear message, rather than inside the MCP client, which reports | ||
| # any failure as a generic connection error. | ||
| auth.get_token() | ||
| transport = StreamableHttpTransport(url, headers=hook.user_agent_header, auth=auth) |
There was a problem hiding this comment.
hook.proxies isn't applied here, so connections that need the proxies extra won't reach the gateway. StreamableHttpTransport takes httpx_client_factory.
| tags: [service] | ||
| - integration-name: Databricks Unity Gateway | ||
| external-doc-url: https://docs.databricks.com/aws/en/unity-gateway/concepts | ||
| how-to-guide: |
There was a problem hiding this comment.
check_doc_files only accepts how-to guides under docs/operators|sensors|transfer, which is why Static checks fails. Drop this how-to-guide (and regenerate get_provider_info.py).
| "apache-airflow-task-sdk", | ||
| "apache-airflow-devel-common", | ||
| "apache-airflow-providers-amazon", | ||
| "apache-airflow-providers-common-ai", |
There was a problem hiding this comment.
apache-airflow-providers-common-ai[mcp]. Otherwise fastmcp isn't installed by uv sync and test_unity_mcp.py skips.
moomindani
left a comment
There was a problem hiding this comment.
One addition to zozo123's review. validate_service_name allows - in each part ([A-Za-z0-9_-]+), and the docs say so. That is not needed to keep the name from adding a path or query to the URL, and Unity Catalog names with hyphens need backtick quoting in every SQL statement, so I would drop - from _SERVICE_NAME_PART and from the docs sentence.
Drafted-by: Claude Code (Opus 5.5); reviewed by @moomindani before posting
Dag authors want their agents to use governed Databricks tools such as Unity Catalog functions, Genie and AI Search through a Unity AI Gateway MCP Service, without putting bearer tokens or gateway URLs in Dag code. The generic MCP toolset needs both in a separate connection and sends a fixed token, which expires during long agent runs with OAuth identities. closes: apache#74193
Each gateway request fetched its token through AirflowToolset.run_blocking, whose lock is shared by every toolset in the process. While another toolset ran a long blocking call, such as a SQL query, every request to the MCP Service waited for it to finish.
Databricks requires USE CATALOG and USE SCHEMA on the parent catalog and schema as well as EXECUTE on the service, so an identity granted only EXECUTE is denied. Users following the guide or the access-denied error would grant too little.
Databricks also requires the calling identity to be assigned to the workspace the request goes to, which the guide did not say.
The MCP client carries on after some gateway error responses, such as one to a notification, but the toolset kept that status and later reported an unrelated failure as it. A tool call with bad arguments then ended the agent run as a gateway error instead of letting the model retry. Also use Databricks' name for the product, Unity Gateway, and correct the guide: Databricks provides MCP Services for workspace tools and SaaS applications, and only external MCP servers can be registered as one.
Tool calls run in parallel, so attributing the last HTTP error seen by the toolset to a failed call could report a server error as a retry the model should make, re-running a call that may already have changed data. A JSON-RPC error the server returns at an HTTP error status, by contrast, is its answer to the call and something the model can correct. Building on common.ai's MCPToolset meant overriding its private internals, which carry no stability guarantee; owning the MCP client keeps the toolset working as common.ai evolves. A missing service is answered with 403, never 404, so it is reported as access denied. Gateway requests now honour the connection's proxies and are only sent over HTTPS, so the bearer token cannot travel in clear text.
48928f5 to
6402af7
Compare
@moomindani thanks, agreed: service names now allow only ASCII letters, digits and |
Without the mcp extra of common.ai, uv sync left out fastmcp and the Unity MCP toolset tests skipped silently. The extra sits in the hand-maintained part of the dev group because the generated part is rewritten without extras.
|
Live test against a Databricks workspace with a PAT, through Unity Gateway MCP Services in front of public MCP servers:
It also found an edge case: Unity Gateway buffers the whole response until the upstream server closes its stream. Servers that answer promptly work fine, but GitMCP keeps its stream open for about 10 s (11 s through the gateway vs 0.7 s direct), which exceeds pydantic-ai's 5 s MCP handshake timeout, so the connection fails with a bare Proposal: add an optional Not tested: a successful tool call through the gateway (blocked by the timeout above), service principal OAuth and Azure AD. |
Dag authors want their agents to use governed Databricks tools such as Unity Catalog functions, Genie and AI Search through a Unity AI Gateway MCP Service, without putting bearer tokens or gateway URLs in Dag code. The generic MCP toolset needs both in a separate connection and sends a fixed token, which expires during long agent runs with OAuth identities.
closes: #74193
Was generative AI tooling used to co-author this PR?
{pr_number}.significant.rst, in airflow-core/newsfragments. You can add this file in a follow-up commit after the PR is created so you know the PR number.