Merge branch 'main' into docs-1.7-release

This commit is contained in:
Mendon Kissling
2025-10-07 17:55:12 -04:00
24 changed files with 1497 additions and 361 deletions

View File

@ -30,7 +30,7 @@ repos:
hooks:
- id: detect-secrets
args: ["--baseline", ".secrets.baseline"]
exclude: ^docs/
exclude: '(^docs/|^SECURITY\.md$)'
- repo: local
hooks:
- id: local-biome-check

View File

@ -14,6 +14,7 @@
[![Ask DeepWiki](https://deepwiki.com/badge.svg)](https://deepwiki.com/langflow-ai/langflow)
> [!CAUTION]
> - Langflow versions 1.6.0 through 1.6.3 have a critical bug where `.env` files are not read, potentially causing security vulnerabilities. **DO NOT** upgrade to these versions if you use `.env` files for configuration. Instead, upgrade to 1.6.4, which includes a fix for this bug.
> - Windows users of Langflow Desktop should **not** use the in-app update feature to upgrade to Langflow version 1.6.0. For upgrade instructions, see [Windows Desktop update issue](https://docs.langflow.org/release-notes#windows-desktop-update-issue).
> - Users must update to Langflow >= 1.3 to protect against [CVE-2025-3248](https://nvd.nist.gov/vuln/detail/CVE-2025-3248)
> - Users must update to Langflow >= 1.5.1 to protect against [CVE-2025-57760](https://github.com/langflow-ai/langflow/security/advisories/GHSA-4gv9-mp8m-592r)

View File

@ -42,6 +42,19 @@ We appreciate your efforts in helping us maintain a secure platform and look for
## Known Vulnerabilities
### Environment Variable Loading Bug (Fixed in 1.6.4)
Langflow versions `1.6.0` through `1.6.3` have a critical bug where environment variables from `.env` files are not being read. This affects all deployments using environment variables for configuration, including security settings.
**Potential security impact:**
- Environment variables from `.env` files are not read.
- Security configurations like `AUTO_LOGIN=false` may not be applied, potentially allowing users to log in as the default superuser.
- Database credentials, API keys, and other sensitive configuration may not be loaded.
**DO NOT** upgrade to Langflow versions `1.6.0` through `1.6.3` if you use `.env` files for configuration. Instead, upgrade to version `1.6.4`, which includes a fix for this bug.
**Fixed in**: Langflow >= 1.6.4
### Code Execution Vulnerability (Fixed in 1.3.0)
Langflow allows users to define and run **custom code components** through endpoints like `/api/v1/validate/code`. In versions < 1.3.0, this endpoint did not enforce authentication or proper sandboxing, allowing **unauthenticated arbitrary code execution**.
@ -99,4 +112,4 @@ export LANGFLOW_SUPERUSER="<your-superuser-username>"
export LANGFLOW_SUPERUSER_PASSWORD="<your-superuser-password>"
export LANGFLOW_DATABASE_URL="<your-production-database-url>" # e.g. "postgresql+psycopg://langflow:secure_pass@db.internal:5432/langflow"
export LANGFLOW_SECRET_KEY="your-strong-random-secret-key"
```
```

View File

@ -122,6 +122,32 @@ You can retrieve flow IDs from the [**API access** pane](/concepts-publish#api-a
Once you have your Langflow server URL, try calling these endpoints that return Langflow metadata.
### Health check
Returns the health status of the Langflow database and chat services:
```bash
curl -X GET \
"$LANGFLOW_SERVER_URL/health_check" \
-H "accept: application/json"
```
<details>
<summary>Result</summary>
```json
{
"status": "ok",
"chat": "ok",
"db": "ok"
}
```
</details>
Langflow provides an additional `GET /health` endpoint.
This endpoint is served by uvicorn before Langflow is fully initialized, so it's not reliable for checking Langflow service health.
### Get version
Returns the current Langflow API version:
@ -130,7 +156,6 @@ Returns the current Langflow API version:
curl -X GET \
"$LANGFLOW_SERVER_URL/api/v1/version" \
-H "accept: application/json"
-H "x-api-key: $LANGFLOW_API_KEY"
```
<details>
@ -138,8 +163,8 @@ curl -X GET \
```text
{
"version": "1.1.1",
"main_version": "1.1.1",
"version": "1.6.0",
"main_version": "1.6.0",
"package": "Langflow"
}
```
@ -154,7 +179,6 @@ Returns configuration details for your Langflow deployment:
curl -X GET \
"$LANGFLOW_SERVER_URL/api/v1/config" \
-H "accept: application/json"
-H "x-api-key: $LANGFLOW_API_KEY"
```
<details>
@ -165,11 +189,21 @@ curl -X GET \
"feature_flags": {
"mvp_components": false
},
"serialization_max_items_length": 1000,
"serialization_max_text_length": 6000,
"frontend_timeout": 0,
"auto_saving": true,
"auto_saving_interval": 1000,
"health_check_max_retries": 5,
"max_file_size_upload": 1024
"max_file_size_upload": 1024,
"webhook_polling_interval": 5000,
"public_flow_cleanup_interval": 3600,
"public_flow_expiration": 86400,
"event_delivery": "streaming",
"webhook_auth_enable": false,
"voice_mode_available": false,
"default_folder_name": "Starter Project",
"hide_getting_started_progress": false
}
```
@ -177,7 +211,8 @@ curl -X GET \
### Get all components
Returns a dictionary of all Langflow components:
Returns a dictionary of all Langflow components.
Requires a [Langflow API key](/api-keys-and-authentication).
```bash
curl -X GET \
@ -215,6 +250,7 @@ Other endpoints are helpful for specific use cases, such as administration and f
* Deployment details:
* GET `/v1/version`: Return Langflow version. See [Get version](/api-reference-api-examples#get-version).
* GET `/v1/config`: Return deployment configuration. See [Get configuration](/api-reference-api-examples#get-configuration).
* GET `/health_check`: Health check endpoint that validates database and chat service connectivity. Returns 500 status if any service is unavailable.
* [Projects endpoints](/api-projects):
* POST `/v1/projects/`: Create a project.

View File

@ -87,6 +87,19 @@ For all changes, see the [Changelog](https://github.com/langflow-ai/langflow/rel
Highlights of this release include the following changes.
For all changes, see the [Changelog](https://github.com/langflow-ai/langflow/releases).
### Known issue, potential security vulnerability: .env file not loaded in versions 1.6.0 through 1.6.3 {#env-file-bug}
Langflow versions 1.6.0 through 1.6.3 have a critical bug where environment variables from `.env` files aren't read.
This affects all deployments using environment variables for configuration, including security settings.
:::warning Potential security vulnerability
If your `.env` file includes `AUTO_LOGIN=false`, upgrading to the impacted versions causes Langflow to fall back to default settings, potentially giving all users superuser access immediately upon upgrade.
Additionally, database credentials, API keys, and other sensitive configurations can't be loaded from `.env` files.
_Don't_ upgrade to any Langflow version from 1.6.0 through 1.6.3 if you use `.env` files for configuration.
Instead, upgrade to 1.6.4, which includes a fix for this bug.
:::
### Known issue: Don't auto-upgrade Windows Desktop {#windows-desktop-update-issue}
:::warning

29
package-lock.json generated Normal file
View File

@ -0,0 +1,29 @@
{
"name": "langflow",
"lockfileVersion": 3,
"requires": true,
"packages": {
"": {
"devDependencies": {
"@types/node": "^24.7.0"
}
},
"node_modules/@types/node": {
"version": "24.7.0",
"resolved": "https://registry.npmjs.org/@types/node/-/node-24.7.0.tgz",
"integrity": "sha512-IbKooQVqUBrlzWTi79E8Fw78l8k1RNtlDDNWsFZs7XonuQSJ8oNYfEeclhprUldXISRMLzBpILuKgPlIxm+/Yw==",
"dev": true,
"license": "MIT",
"dependencies": {
"undici-types": "~7.14.0"
}
},
"node_modules/undici-types": {
"version": "7.14.0",
"resolved": "https://registry.npmjs.org/undici-types/-/undici-types-7.14.0.tgz",
"integrity": "sha512-QQiYxHuyZ9gQUIrmPo3IA+hUl4KYk8uSA7cHrcKd/l3p1OTpZcM0Tbp9x7FAtXdAYhlasd60ncPpgu6ihG6TOA==",
"dev": true,
"license": "MIT"
}
}
}

5
package.json Normal file
View File

@ -0,0 +1,5 @@
{
"devDependencies": {
"@types/node": "^24.7.0"
}
}

File diff suppressed because one or more lines are too long

View File

@ -147,13 +147,19 @@ async def aupdate_messages(messages: Message | list[Message]) -> list[Message]:
if msg.flow_id and isinstance(msg.flow_id, str):
msg.flow_id = UUID(msg.flow_id)
session.add(msg)
await session.commit()
await session.refresh(msg)
updated_messages.append(msg)
else:
error_message = f"Message with id {message.id} not found"
await logger.awarning(error_message)
raise ValueError(error_message)
# Batch commit all updates at once
await session.commit()
# Skip refresh during commit - the msg objects already have the updated values
# Refresh is only needed if we need database-generated values (like timestamps)
# For streaming performance, we skip this extra round-trip
return [MessageRead.model_validate(message, from_attributes=True) for message in updated_messages]

View File

@ -13,7 +13,14 @@ from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from lfx.base.mcp import util
from lfx.base.mcp.util import MCPSessionManager, MCPSseClient, MCPStdioClient, _process_headers, validate_headers
from lfx.base.mcp.util import (
MCPSessionManager,
MCPSseClient,
MCPStdioClient,
MCPStreamableHttpClient,
_process_headers,
validate_headers,
)
class TestMCPSessionManager:
@ -186,11 +193,11 @@ class TestHeaderValidation:
assert result == {"safe-header": "safe-value"}
class TestSSEHeaderIntegration:
"""Integration test to verify headers are properly passed through the entire SSE flow."""
class TestStreamableHTTPHeaderIntegration:
"""Integration test to verify headers are properly passed through the entire StreamableHTTP flow."""
async def test_headers_processing(self):
"""Test that headers flow properly from server config through to SSE client connection."""
"""Test that headers flow properly from server config through to StreamableHTTP client connection."""
# Test the header processing function directly
headers_input = [
{"key": "Authorization", "value": "Bearer test-token"},
@ -231,15 +238,15 @@ class TestSSEHeaderIntegration:
result = _process_headers(invalid_headers)
assert result == {"valid-header": "good"}
async def test_sse_client_header_storage(self):
async def test_streamable_http_client_header_storage(self):
"""Test that SSE client properly stores headers in connection params."""
sse_client = MCPSseClient()
streamable_http_client = MCPStreamableHttpClient()
test_url = "http://test.url"
test_headers = {"Authorization": "Bearer test123", "Custom": "value"}
# Test that headers are properly stored in connection params
# Set connection params as a dict like the implementation expects
sse_client._connection_params = {
streamable_http_client._connection_params = {
"url": test_url,
"headers": test_headers,
"timeout_seconds": 30,
@ -247,8 +254,8 @@ class TestSSEHeaderIntegration:
}
# Verify headers are stored
assert sse_client._connection_params["url"] == test_url
assert sse_client._connection_params["headers"] == test_headers
assert streamable_http_client._connection_params["url"] == test_url
assert streamable_http_client._connection_params["headers"] == test_headers
class TestFieldNameConversion:
@ -930,22 +937,22 @@ class TestMCPStdioClientWithEverythingServer:
await stdio_client.disconnect()
class TestMCPSseClientWithDeepWikiServer:
class TestMCPStreamableHttpClientWithDeepWikiServer:
"""Test MCPSseClient with the DeepWiki MCP server."""
@pytest.fixture
def sse_client(self):
def streamable_http_client(self):
"""Create an SSE client for testing."""
return MCPSseClient()
return MCPStreamableHttpClient()
@pytest.mark.asyncio
async def test_connect_to_deepwiki_server(self, sse_client):
async def test_connect_to_deepwiki_server(self, streamable_http_client):
"""Test connecting to the DeepWiki MCP server."""
url = "https://mcp.deepwiki.com/sse"
try:
# Connect to the server
tools = await sse_client.connect_to_server(url)
tools = await streamable_http_client.connect_to_server(url)
# Verify tools were returned
assert len(tools) > 0
@ -961,16 +968,16 @@ class TestMCPSseClientWithDeepWikiServer:
# If the server is not accessible, skip the test
pytest.skip(f"DeepWiki server not accessible: {e}")
finally:
await sse_client.disconnect()
await streamable_http_client.disconnect()
@pytest.mark.asyncio
async def test_run_wiki_structure_tool(self, sse_client):
async def test_run_wiki_structure_tool(self, streamable_http_client):
"""Test running the read_wiki_structure tool."""
url = "https://mcp.deepwiki.com/sse"
try:
# Connect to the server
tools = await sse_client.connect_to_server(url)
tools = await streamable_http_client.connect_to_server(url)
# Find the read_wiki_structure tool
wiki_tool = None
@ -982,7 +989,7 @@ class TestMCPSseClientWithDeepWikiServer:
assert wiki_tool is not None, "read_wiki_structure tool not found"
# Run the tool with a test repository (use repoName as expected by the API)
result = await sse_client.run_tool("read_wiki_structure", {"repoName": "microsoft/vscode"})
result = await streamable_http_client.run_tool("read_wiki_structure", {"repoName": "microsoft/vscode"})
# Verify the result
assert result is not None
@ -993,16 +1000,16 @@ class TestMCPSseClientWithDeepWikiServer:
# If the server is not accessible or the tool fails, skip the test
pytest.skip(f"DeepWiki server test failed: {e}")
finally:
await sse_client.disconnect()
await streamable_http_client.disconnect()
@pytest.mark.asyncio
async def test_ask_question_tool(self, sse_client):
async def test_ask_question_tool(self, streamable_http_client):
"""Test running the ask_question tool."""
url = "https://mcp.deepwiki.com/sse"
try:
# Connect to the server
tools = await sse_client.connect_to_server(url)
tools = await streamable_http_client.connect_to_server(url)
# Find the ask_question tool
ask_tool = None
@ -1014,7 +1021,7 @@ class TestMCPSseClientWithDeepWikiServer:
assert ask_tool is not None, "ask_question tool not found"
# Run the tool with a test question (use repoName as expected by the API)
result = await sse_client.run_tool(
result = await streamable_http_client.run_tool(
"ask_question", {"repoName": "microsoft/vscode", "question": "What is VS Code?"}
)
@ -1027,14 +1034,14 @@ class TestMCPSseClientWithDeepWikiServer:
# If the server is not accessible or the tool fails, skip the test
pytest.skip(f"DeepWiki server test failed: {e}")
finally:
await sse_client.disconnect()
await streamable_http_client.disconnect()
@pytest.mark.asyncio
async def test_url_validation(self, sse_client):
async def test_url_validation(self, streamable_http_client):
"""Test URL validation for SSE connections."""
# Test valid URL
valid_url = "https://mcp.deepwiki.com/sse"
is_valid, error = await sse_client.validate_url(valid_url)
is_valid, error = await streamable_http_client.validate_url(valid_url)
# Either valid or accessible, or rate-limited (429) which indicates server is reachable
if not is_valid and "429" in error:
# Rate limiting indicates the server is accessible but limiting requests
@ -1044,29 +1051,10 @@ class TestMCPSseClientWithDeepWikiServer:
# Test invalid URL
invalid_url = "not_a_url"
is_valid, error = await sse_client.validate_url(invalid_url)
is_valid, error = await streamable_http_client.validate_url(invalid_url)
assert not is_valid
assert error != ""
@pytest.mark.asyncio
async def test_redirect_handling(self, sse_client):
"""Test redirect handling for SSE connections."""
# Test with the DeepWiki URL
url = "https://mcp.deepwiki.com/sse"
try:
# Check for redirects
final_url = await sse_client.pre_check_redirect(url)
# Should return a URL (either original or redirected)
assert final_url is not None
assert isinstance(final_url, str)
assert final_url.startswith("http")
except Exception as e:
# If the server is not accessible, skip the test
pytest.skip(f"DeepWiki server not accessible for redirect test: {e}")
@pytest.fixture
def mock_tool(self):
"""Create a mock MCP tool."""
@ -1117,14 +1105,14 @@ class TestMCPSseClientUnit:
mock_response.status_code = 200
mock_client.return_value.__aenter__.return_value.get.return_value = mock_response
is_valid, error_msg = await sse_client.validate_url("http://test.url", {})
is_valid, error_msg = await sse_client.validate_url("http://test.url")
assert is_valid is True
assert error_msg == ""
async def test_validate_url_invalid_format(self, sse_client):
"""Test URL validation with invalid format."""
is_valid, error_msg = await sse_client.validate_url("invalid-url", {})
is_valid, error_msg = await sse_client.validate_url("invalid-url")
assert is_valid is False
assert "Invalid URL format" in error_msg
@ -1136,7 +1124,7 @@ class TestMCPSseClientUnit:
mock_response.status_code = 404
mock_client.return_value.__aenter__.return_value.get.return_value = mock_response
is_valid, error_msg = await sse_client.validate_url("http://test.url", {})
is_valid, error_msg = await sse_client.validate_url("http://test.url")
assert is_valid is True
assert error_msg == ""
@ -1149,7 +1137,6 @@ class TestMCPSseClientUnit:
with (
patch.object(sse_client, "validate_url", return_value=(True, "")),
patch.object(sse_client, "pre_check_redirect", return_value=test_url),
patch.object(sse_client, "_get_or_create_session") as mock_get_session,
):
# Mock session
@ -1194,29 +1181,10 @@ class TestMCPSseClientUnit:
result_session = await sse_client._get_or_create_session()
# Verify session manager was called with correct parameters including normalized headers
mock_manager.get_session.assert_called_once_with("test_context", sse_client._connection_params, "sse")
assert result_session == mock_session
async def test_pre_check_redirect_with_headers(self, sse_client):
"""Test pre-check redirect functionality with custom headers."""
test_url = "http://test.url"
redirect_url = "http://redirect.url"
# Use pre-validated headers since pre_check_redirect expects already validated headers
test_headers = {"authorization": "Bearer token123"} # already normalized
with patch("httpx.AsyncClient") as mock_client:
mock_response = MagicMock()
mock_response.status_code = 307
mock_response.headers.get.return_value = redirect_url
mock_client.return_value.__aenter__.return_value.get.return_value = mock_response
result = await sse_client.pre_check_redirect(test_url, test_headers)
assert result == redirect_url
# Verify validated headers were passed to the request
mock_client.return_value.__aenter__.return_value.get.assert_called_with(
test_url, timeout=2.0, headers={"Accept": "text/event-stream", **test_headers}
mock_manager.get_session.assert_called_once_with(
"test_context", sse_client._connection_params, "streamable_http"
)
assert result_session == mock_session
async def test_run_tool_with_retry_on_connection_error(self, sse_client):
"""Test that run_tool retries on connection errors."""

View File

@ -27,7 +27,7 @@ async def create_event_iterator(events: list[dict[str, Any]]) -> AsyncIterator[d
async def test_chain_start_event():
"""Test handling of on_chain_start event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
events = [
{"event": "on_chain_start", "data": {"input": {"input": "test input", "chat_history": []}}, "start_time": 0}
@ -52,7 +52,7 @@ async def test_chain_start_event():
async def test_chain_end_event():
"""Test handling of on_chain_end event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
# Create a mock AgentFinish output
output = AgentFinish(return_values={"output": "final output"}, log="test log")
@ -81,7 +81,7 @@ async def test_tool_start_event():
send_message = AsyncMock()
# Set up the send_message mock to return the modified message
def update_message(message):
def update_message(message, skip_db_update=False): # noqa: ARG001, FBT002
# Return a copy of the message to simulate real behavior
return Message(**message.model_dump())
@ -117,7 +117,7 @@ async def test_tool_start_event():
async def test_tool_end_event():
"""Test handling of on_tool_end event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
events = [
{
@ -152,7 +152,7 @@ async def test_tool_end_event():
async def test_tool_error_event():
"""Test handling of on_tool_error event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
events = [
{
@ -188,7 +188,7 @@ async def test_tool_error_event():
async def test_chain_stream_event():
"""Test handling of on_chain_stream event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
events = [{"event": "on_chain_stream", "data": {"chunk": {"output": "streamed output"}}, "start_time": 0}]
agent_message = Message(
@ -206,7 +206,7 @@ async def test_chain_stream_event():
async def test_multiple_events():
"""Test handling of multiple events in sequence."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
# Create a mock AgentFinish output instead of MockOutput
output = AgentFinish(return_values={"output": "final output"}, log="test log")
@ -249,7 +249,7 @@ async def test_multiple_events():
async def test_unknown_event():
"""Test handling of unknown event type."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -274,7 +274,7 @@ async def test_unknown_event():
async def test_handle_on_chain_start_with_input():
"""Test handle_on_chain_start with input."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -293,7 +293,7 @@ async def test_handle_on_chain_start_with_input():
async def test_handle_on_chain_start_no_input():
"""Test handle_on_chain_start without input."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -312,7 +312,7 @@ async def test_handle_on_chain_start_no_input():
async def test_handle_on_chain_end_with_output():
"""Test handle_on_chain_end with output."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -333,7 +333,7 @@ async def test_handle_on_chain_end_with_output():
async def test_handle_on_chain_end_no_output():
"""Test handle_on_chain_end without output key in data."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -352,7 +352,7 @@ async def test_handle_on_chain_end_no_output():
async def test_handle_on_chain_end_empty_data():
"""Test handle_on_chain_end with empty data."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -371,7 +371,7 @@ async def test_handle_on_chain_end_empty_data():
async def test_handle_on_chain_end_with_empty_return_values():
"""Test handle_on_chain_end with empty return_values."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -395,7 +395,7 @@ async def test_handle_on_chain_end_with_empty_return_values():
async def test_handle_on_tool_start():
"""Test handle_on_tool_start event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
tool_blocks_map = {}
agent_message = Message(
sender=MESSAGE_SENDER_AI,
@ -427,7 +427,7 @@ async def test_handle_on_tool_start():
async def test_handle_on_tool_end():
"""Test handle_on_tool_end event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
tool_blocks_map = {}
agent_message = Message(
sender=MESSAGE_SENDER_AI,
@ -464,7 +464,7 @@ async def test_handle_on_tool_end():
async def test_handle_on_tool_error():
"""Test handle_on_tool_error event."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
tool_blocks_map = {}
agent_message = Message(
sender=MESSAGE_SENDER_AI,
@ -503,7 +503,7 @@ async def test_handle_on_tool_error():
async def test_handle_on_chain_stream_with_output():
"""Test handle_on_chain_stream with output."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",
@ -524,7 +524,7 @@ async def test_handle_on_chain_stream_with_output():
async def test_handle_on_chain_stream_no_output():
"""Test handle_on_chain_stream without output."""
send_message = AsyncMock(side_effect=lambda message: message)
send_message = AsyncMock(side_effect=lambda message, skip_db_update=False: message) # noqa: ARG005
agent_message = Message(
sender=MESSAGE_SENDER_AI,
sender_name="Agent",

View File

@ -2,14 +2,14 @@
This test suite validates the MCP component functionality using real MCP servers:
- Everything server (stdio mode) - provides echo and other tools
- DeepWiki server (SSE mode) - provides wiki-related tools
- HTTP/SSE servers (streamable HTTP mode) - provides various tools
"""
import shutil
from unittest.mock import AsyncMock, MagicMock, patch
import pytest
from lfx.base.mcp.util import MCPSessionManager, MCPSseClient, MCPStdioClient
from lfx.base.mcp.util import MCPSessionManager, MCPStdioClient, MCPStreamableHttpClient
from lfx.components.agents.mcp_component import MCPToolsComponent
from tests.base import ComponentTestBaseWithoutClient, VersionComponentMapping
@ -45,9 +45,9 @@ class TestMCPToolsComponent(ComponentTestBaseWithoutClient):
# Check that the component has the expected attributes
assert hasattr(component, "stdio_client")
assert hasattr(component, "sse_client")
assert hasattr(component, "streamable_http_client")
assert isinstance(component.stdio_client, MCPStdioClient)
assert isinstance(component.sse_client, MCPSseClient)
assert isinstance(component.streamable_http_client, MCPStreamableHttpClient)
# Check that the component has a session manager
session_manager = component.stdio_client._get_session_manager()
@ -87,11 +87,11 @@ class TestMCPToolsComponentIntegration:
pytest.skip(f"Everything server not accessible: {e}")
@pytest.mark.asyncio
async def test_sse_mode_integration(self, component):
"""Test the component in SSE mode with DeepWiki server."""
# Configure for SSE mode
component.mode = "SSE"
component.sse_url = "https://mcp.deepwiki.com/sse"
async def test_streamable_http_mode_integration(self, component):
"""Test the component in Streamable HTTP mode with DeepWiki server."""
# Configure for Streamable HTTP mode
component.mode = "Streamable HTTP"
component.streamable_http_url = "https://mcp.deepwiki.com/mcp"
try:
# Mock the update_tool_list method to simulate server connection
@ -111,27 +111,27 @@ class TestMCPToolsComponentIntegration:
@pytest.mark.asyncio
async def test_session_context_setting(self, component):
"""Test that session context is properly set."""
# Set session context
# Set session context on both clients
component.stdio_client.set_session_context("test_context")
component.sse_client.set_session_context("test_context")
component.streamable_http_client.set_session_context("test_context")
# Verify context was set
assert component.stdio_client._session_context == "test_context"
assert component.sse_client._session_context == "test_context"
assert component.streamable_http_client._session_context == "test_context"
@pytest.mark.asyncio
async def test_session_manager_sharing(self, component):
"""Test that session managers are shared through component cache."""
# Get session managers
# Get session managers from both clients
stdio_manager = component.stdio_client._get_session_manager()
sse_manager = component.sse_client._get_session_manager()
http_manager = component.streamable_http_client._get_session_manager()
# Both should be MCPSessionManager instances
assert isinstance(stdio_manager, MCPSessionManager)
assert isinstance(sse_manager, MCPSessionManager)
assert isinstance(http_manager, MCPSessionManager)
# They should be the same instance (shared through cache)
assert stdio_manager is sse_manager
assert stdio_manager is http_manager
class TestMCPComponentErrorHandling:

View File

@ -98,6 +98,7 @@
"devDependencies": {
"@biomejs/biome": "2.1.1",
"@jest/types": "^30.0.1",
"@modelcontextprotocol/server-everything": "^0.6.2",
"@playwright/test": "^1.52.0",
"@swc/cli": "^0.5.2",
"@swc/core": "^1.6.1",
@ -3898,6 +3899,34 @@
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/@modelcontextprotocol/sdk": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/@modelcontextprotocol/sdk/-/sdk-1.0.1.tgz",
"integrity": "sha512-slLdFaxQJ9AlRg+hw28iiTtGvShAOgOKXcD0F91nUcRYiOMuS9ZBYjcdNZRXW9G5JQ511GRTdUy1zQVZDpJ+4w==",
"dev": true,
"license": "MIT",
"dependencies": {
"content-type": "^1.0.5",
"raw-body": "^3.0.0",
"zod": "^3.23.8"
}
},
"node_modules/@modelcontextprotocol/server-everything": {
"version": "0.6.2",
"resolved": "https://registry.npmjs.org/@modelcontextprotocol/server-everything/-/server-everything-0.6.2.tgz",
"integrity": "sha512-8ILXxbM8kBWbrtEZoCBYqvAPRyPHhayo4aS2llIPt9oi55o4EUj/90xTFeXuyiM3ReQR/fov3ZmpCy/vC+svUQ==",
"dev": true,
"license": "MIT",
"dependencies": {
"@modelcontextprotocol/sdk": "1.0.1",
"express": "^4.21.1",
"zod": "^3.23.8",
"zod-to-json-schema": "^3.23.5"
},
"bin": {
"mcp-server-everything": "dist/index.js"
}
},
"node_modules/@napi-rs/nice": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/@napi-rs/nice/-/nice-1.0.1.tgz",
@ -8307,6 +8336,13 @@
"dequal": "^2.0.3"
}
},
"node_modules/array-flatten": {
"version": "1.1.1",
"resolved": "https://registry.npmjs.org/array-flatten/-/array-flatten-1.1.1.tgz",
"integrity": "sha512-PCVAQswWemu6UdxsDFFX/+gVeYqKAod3D3UVm91jHwynguOwAvYPhx8nNlM++NqRcK6CxxpUafjmhIdKiHibqg==",
"dev": true,
"license": "MIT"
},
"node_modules/ast-types": {
"version": "0.14.2",
"resolved": "https://registry.npmjs.org/ast-types/-/ast-types-0.14.2.tgz",
@ -8635,6 +8671,77 @@
"readable-stream": "^3.4.0"
}
},
"node_modules/body-parser": {
"version": "1.20.3",
"resolved": "https://registry.npmjs.org/body-parser/-/body-parser-1.20.3.tgz",
"integrity": "sha512-7rAxByjUMqQ3/bHJy7D6OGXvx/MMc4IqBn/X0fcM1QUcAItpZrBEYhWGem+tzXH90c+G01ypMcYJBO9Y30203g==",
"dev": true,
"license": "MIT",
"dependencies": {
"bytes": "3.1.2",
"content-type": "~1.0.5",
"debug": "2.6.9",
"depd": "2.0.0",
"destroy": "1.2.0",
"http-errors": "2.0.0",
"iconv-lite": "0.4.24",
"on-finished": "2.4.1",
"qs": "6.13.0",
"raw-body": "2.5.2",
"type-is": "~1.6.18",
"unpipe": "1.0.0"
},
"engines": {
"node": ">= 0.8",
"npm": "1.2.8000 || >= 1.4.16"
}
},
"node_modules/body-parser/node_modules/debug": {
"version": "2.6.9",
"resolved": "https://registry.npmjs.org/debug/-/debug-2.6.9.tgz",
"integrity": "sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA==",
"dev": true,
"license": "MIT",
"dependencies": {
"ms": "2.0.0"
}
},
"node_modules/body-parser/node_modules/iconv-lite": {
"version": "0.4.24",
"resolved": "https://registry.npmjs.org/iconv-lite/-/iconv-lite-0.4.24.tgz",
"integrity": "sha512-v3MXnZAcvnywkTUEZomIActle7RXXeedOR31wwl7VlyoXO4Qi9arvSenNQWne1TcRwhCL1HwLI21bEqdpj8/rA==",
"dev": true,
"license": "MIT",
"dependencies": {
"safer-buffer": ">= 2.1.2 < 3"
},
"engines": {
"node": ">=0.10.0"
}
},
"node_modules/body-parser/node_modules/ms": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz",
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
"dev": true,
"license": "MIT"
},
"node_modules/body-parser/node_modules/raw-body": {
"version": "2.5.2",
"resolved": "https://registry.npmjs.org/raw-body/-/raw-body-2.5.2.tgz",
"integrity": "sha512-8zGqypfENjCIqGhgXToC8aB2r7YrBX+AQAfIPs/Mlk+BtPTztOvTS01NRW/3Eh60J+a48lt8qsCzirQ6loCVfA==",
"dev": true,
"license": "MIT",
"dependencies": {
"bytes": "3.1.2",
"http-errors": "2.0.0",
"iconv-lite": "0.4.24",
"unpipe": "1.0.0"
},
"engines": {
"node": ">= 0.8"
}
},
"node_modules/boxen": {
"version": "5.1.2",
"resolved": "https://registry.npmjs.org/boxen/-/boxen-5.1.2.tgz",
@ -8821,6 +8928,16 @@
"integrity": "sha512-E+XQCRwSbaaiChtv6k6Dwgc+bx+Bs6vuKJHHl5kox/BaKbhiXzqQOwK4cO22yElGp2OCmjwVhT3HmxgyPGnJfQ==",
"dev": true
},
"node_modules/bytes": {
"version": "3.1.2",
"resolved": "https://registry.npmjs.org/bytes/-/bytes-3.1.2.tgz",
"integrity": "sha512-/Nf7TyzTx6S3yRJObOAV7956r8cr2+Oj8AC5dt8wSP3BQAoeX58NoHyCU8P8zGkNXStjTSi6fzO6F0pBdcYbEg==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/cacheable-lookup": {
"version": "7.0.0",
"resolved": "https://registry.npmjs.org/cacheable-lookup/-/cacheable-lookup-7.0.0.tgz",
@ -8863,6 +8980,23 @@
"node": ">= 0.4"
}
},
"node_modules/call-bound": {
"version": "1.0.4",
"resolved": "https://registry.npmjs.org/call-bound/-/call-bound-1.0.4.tgz",
"integrity": "sha512-+ys997U96po4Kx/ABpBCqhA9EuxJaQWDQg7295H4hBphv3IZg0boBKuwYpt4YXp6MZ5AmZQnU/tyMTlRpaSejg==",
"dev": true,
"license": "MIT",
"dependencies": {
"call-bind-apply-helpers": "^1.0.2",
"get-intrinsic": "^1.3.0"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/callsites": {
"version": "3.1.0",
"resolved": "https://registry.npmjs.org/callsites/-/callsites-3.1.0.tgz",
@ -9372,6 +9506,16 @@
"node": ">= 0.6"
}
},
"node_modules/content-type": {
"version": "1.0.5",
"resolved": "https://registry.npmjs.org/content-type/-/content-type-1.0.5.tgz",
"integrity": "sha512-nTjqfcBFEipKdXCv4YDQWCfmcLZKm81ldF0pAopTvyrFGVbcR6P/VAAd5G7N+0tTr8QqiU0tFadD6FK4NtJwOA==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/convert-source-map": {
"version": "1.9.0",
"resolved": "https://registry.npmjs.org/convert-source-map/-/convert-source-map-1.9.0.tgz",
@ -9388,6 +9532,13 @@
"node": ">= 0.6"
}
},
"node_modules/cookie-signature": {
"version": "1.0.6",
"resolved": "https://registry.npmjs.org/cookie-signature/-/cookie-signature-1.0.6.tgz",
"integrity": "sha512-QADzlaHc8icV8I7vbaJXJwod9HWYp8uCqf1xa4OfNu1T7JVxQIrUgOWtHdNDtPiywmFbiS12VjotIXLrKM3orQ==",
"dev": true,
"license": "MIT"
},
"node_modules/cors": {
"version": "2.8.5",
"resolved": "https://registry.npmjs.org/cors/-/cors-2.8.5.tgz",
@ -9806,6 +9957,16 @@
"optional": true,
"peer": true
},
"node_modules/depd": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/depd/-/depd-2.0.0.tgz",
"integrity": "sha512-g7nH6P6dyDioJogAAGprGpCtVImJhpPk/roCzdb3fIh61/s/nPsfR6onyMwkCAR/OlC3yBC0lESvUoQEAssIrw==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/dequal": {
"version": "2.0.3",
"resolved": "https://registry.npmjs.org/dequal/-/dequal-2.0.3.tgz",
@ -9815,6 +9976,17 @@
"node": ">=6"
}
},
"node_modules/destroy": {
"version": "1.2.0",
"resolved": "https://registry.npmjs.org/destroy/-/destroy-1.2.0.tgz",
"integrity": "sha512-2sJGJTaXIIaR1w4iJSNoN0hnMY7Gpc/n8D4qSCJw8QqFWXf7cuAgnEHxBpweaVcPevC2l3KpjYCx3NypQQgaJg==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8",
"npm": "1.2.8000 || >= 1.4.16"
}
},
"node_modules/detect-libc": {
"version": "2.0.4",
"resolved": "https://registry.npmjs.org/detect-libc/-/detect-libc-2.0.4.tgz",
@ -9982,6 +10154,13 @@
"integrity": "sha512-I88TYZWc9XiYHRQ4/3c5rjjfgkjhLyW2luGIheGERbNQ6OY7yTybanSpDXZa8y7VUP9YmDcYa+eyq4ca7iLqWA==",
"license": "MIT"
},
"node_modules/ee-first": {
"version": "1.1.1",
"resolved": "https://registry.npmjs.org/ee-first/-/ee-first-1.1.1.tgz",
"integrity": "sha512-WMwm9LhRUo+WUaRN+vRuETqG89IgZphVSNkdFgeb6sS/E4OrDIN7t48CAewSHXc6C8lefD8KKfr5vY61brQlow==",
"dev": true,
"license": "MIT"
},
"node_modules/effect": {
"version": "3.16.4",
"resolved": "https://registry.npmjs.org/effect/-/effect-3.16.4.tgz",
@ -10037,6 +10216,16 @@
"integrity": "sha512-EC+0oUMY1Rqm4O6LLrgjtYDvcVYTy7chDnM4Q7030tP4Kwj3u/pR6gP9ygnp2CJMK5Gq+9Q2oqmrFJAz01DXjw==",
"license": "MIT"
},
"node_modules/encodeurl": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/encodeurl/-/encodeurl-2.0.0.tgz",
"integrity": "sha512-Q0n9HRi4m6JuGIV1eFlmvJB7ZEVxu93IrMyiMsGC0lrMJMWzRgx6WGquyfQgZVb31vhGgXnfmPNNXmxnOkRBrg==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/end-of-stream": {
"version": "1.4.4",
"resolved": "https://registry.npmjs.org/end-of-stream/-/end-of-stream-1.4.4.tgz",
@ -10244,6 +10433,13 @@
"node": ">=8"
}
},
"node_modules/escape-html": {
"version": "1.0.3",
"resolved": "https://registry.npmjs.org/escape-html/-/escape-html-1.0.3.tgz",
"integrity": "sha512-NiSupZ4OeuGwr68lGIeym/ksIZMJodUGOSCZ/FSnTxcrekbvqrgdUxlJOMpijaKZVjAJrWrGs/6Jy8OMuyj9ow==",
"dev": true,
"license": "MIT"
},
"node_modules/escape-string-regexp": {
"version": "4.0.0",
"resolved": "https://registry.npmjs.org/escape-string-regexp/-/escape-string-regexp-4.0.0.tgz",
@ -10349,6 +10545,16 @@
"node": ">=0.10.0"
}
},
"node_modules/etag": {
"version": "1.8.1",
"resolved": "https://registry.npmjs.org/etag/-/etag-1.8.1.tgz",
"integrity": "sha512-aIL5Fx7mawVa300al2BnEE4iNvo1qETxLrPI/o05L7z6go7fCw1J6EQmbK4FmJ2AS7kgVF/KEZWufBfdClMcPg==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/execa": {
"version": "5.1.1",
"resolved": "https://registry.npmjs.org/execa/-/execa-5.1.1.tgz",
@ -10409,6 +10615,80 @@
"node": "^14.15.0 || ^16.10.0 || >=18.0.0"
}
},
"node_modules/express": {
"version": "4.21.2",
"resolved": "https://registry.npmjs.org/express/-/express-4.21.2.tgz",
"integrity": "sha512-28HqgMZAmih1Czt9ny7qr6ek2qddF4FclbMzwhCREB6OFfH+rXAnuNCwo1/wFvrtbgsQDb4kSbX9de9lFbrXnA==",
"dev": true,
"license": "MIT",
"dependencies": {
"accepts": "~1.3.8",
"array-flatten": "1.1.1",
"body-parser": "1.20.3",
"content-disposition": "0.5.4",
"content-type": "~1.0.4",
"cookie": "0.7.1",
"cookie-signature": "1.0.6",
"debug": "2.6.9",
"depd": "2.0.0",
"encodeurl": "~2.0.0",
"escape-html": "~1.0.3",
"etag": "~1.8.1",
"finalhandler": "1.3.1",
"fresh": "0.5.2",
"http-errors": "2.0.0",
"merge-descriptors": "1.0.3",
"methods": "~1.1.2",
"on-finished": "2.4.1",
"parseurl": "~1.3.3",
"path-to-regexp": "0.1.12",
"proxy-addr": "~2.0.7",
"qs": "6.13.0",
"range-parser": "~1.2.1",
"safe-buffer": "5.2.1",
"send": "0.19.0",
"serve-static": "1.16.2",
"setprototypeof": "1.2.0",
"statuses": "2.0.1",
"type-is": "~1.6.18",
"utils-merge": "1.0.1",
"vary": "~1.1.2"
},
"engines": {
"node": ">= 0.10.0"
},
"funding": {
"type": "opencollective",
"url": "https://opencollective.com/express"
}
},
"node_modules/express/node_modules/cookie": {
"version": "0.7.1",
"resolved": "https://registry.npmjs.org/cookie/-/cookie-0.7.1.tgz",
"integrity": "sha512-6DnInpx7SJ2AK3+CTUE/ZM0vWTUboZCegxhC2xiIydHR9jNuTAASBrfEpHhiGOZw/nX51bHt6YQl8jsGo4y/0w==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/express/node_modules/debug": {
"version": "2.6.9",
"resolved": "https://registry.npmjs.org/debug/-/debug-2.6.9.tgz",
"integrity": "sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA==",
"dev": true,
"license": "MIT",
"dependencies": {
"ms": "2.0.0"
}
},
"node_modules/express/node_modules/ms": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz",
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
"dev": true,
"license": "MIT"
},
"node_modules/ext-list": {
"version": "2.2.2",
"resolved": "https://registry.npmjs.org/ext-list/-/ext-list-2.2.2.tgz",
@ -10706,6 +10986,42 @@
"node": ">=8"
}
},
"node_modules/finalhandler": {
"version": "1.3.1",
"resolved": "https://registry.npmjs.org/finalhandler/-/finalhandler-1.3.1.tgz",
"integrity": "sha512-6BN9trH7bp3qvnrRyzsBz+g3lZxTNZTbVO2EV1CS0WIcDbawYVdYvGflME/9QP0h0pYlCDBCTjYa9nZzMDpyxQ==",
"dev": true,
"license": "MIT",
"dependencies": {
"debug": "2.6.9",
"encodeurl": "~2.0.0",
"escape-html": "~1.0.3",
"on-finished": "2.4.1",
"parseurl": "~1.3.3",
"statuses": "2.0.1",
"unpipe": "~1.0.0"
},
"engines": {
"node": ">= 0.8"
}
},
"node_modules/finalhandler/node_modules/debug": {
"version": "2.6.9",
"resolved": "https://registry.npmjs.org/debug/-/debug-2.6.9.tgz",
"integrity": "sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA==",
"dev": true,
"license": "MIT",
"dependencies": {
"ms": "2.0.0"
}
},
"node_modules/finalhandler/node_modules/ms": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz",
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
"dev": true,
"license": "MIT"
},
"node_modules/find-root": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/find-root/-/find-root-1.1.0.tgz",
@ -10811,6 +11127,16 @@
"node": ">=0.4.x"
}
},
"node_modules/forwarded": {
"version": "0.2.0",
"resolved": "https://registry.npmjs.org/forwarded/-/forwarded-0.2.0.tgz",
"integrity": "sha512-buRG0fpBtRHSTCOASe6hD258tEubFoRLb4ZNA6NxMVHNw2gOcwHo9wyablzMzOA5z9xA9L1KNjk/Nt6MT9aYow==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/fraction.js": {
"version": "4.3.7",
"resolved": "https://registry.npmjs.org/fraction.js/-/fraction.js-4.3.7.tgz",
@ -10869,6 +11195,16 @@
"license": "0BSD",
"peer": true
},
"node_modules/fresh": {
"version": "0.5.2",
"resolved": "https://registry.npmjs.org/fresh/-/fresh-0.5.2.tgz",
"integrity": "sha512-zJ2mQYM18rEFOudeV4GShTGIQ7RbzA7ozbU9I/XBpm7kqgMywgmylMwXHxZJmkVoYkna9d2pVXVXPdYTP9ej8Q==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/fs-constants": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/fs-constants/-/fs-constants-1.0.0.tgz",
@ -11602,6 +11938,23 @@
"dev": true,
"license": "BSD-2-Clause"
},
"node_modules/http-errors": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/http-errors/-/http-errors-2.0.0.tgz",
"integrity": "sha512-FtwrG/euBzaEjYeRqOgly7G0qviiXoJWnvEH2Z1plBdXgbyjv34pHTSb9zoeHMyDy33+DWy5Wt9Wo+TURtOYSQ==",
"dev": true,
"license": "MIT",
"dependencies": {
"depd": "2.0.0",
"inherits": "2.0.4",
"setprototypeof": "1.2.0",
"statuses": "2.0.1",
"toidentifier": "1.0.1"
},
"engines": {
"node": ">= 0.8"
}
},
"node_modules/http-proxy-agent": {
"version": "5.0.0",
"resolved": "https://registry.npmjs.org/http-proxy-agent/-/http-proxy-agent-5.0.0.tgz",
@ -11796,6 +12149,16 @@
"kind-of": "^6.0.2"
}
},
"node_modules/ipaddr.js": {
"version": "1.9.1",
"resolved": "https://registry.npmjs.org/ipaddr.js/-/ipaddr.js-1.9.1.tgz",
"integrity": "sha512-0KI/607xoxSToH7GjN1FfSbLoU0+btTicjsQSWQlh/hZykN8KpmMf7uYwPW3R+akZ6R/w18ZlXSHBYXiYUPO3g==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.10"
}
},
"node_modules/is-alphabetical": {
"version": "1.0.4",
"resolved": "https://registry.npmjs.org/is-alphabetical/-/is-alphabetical-1.0.4.tgz",
@ -15422,12 +15785,32 @@
"url": "https://opencollective.com/unified"
}
},
"node_modules/media-typer": {
"version": "0.3.0",
"resolved": "https://registry.npmjs.org/media-typer/-/media-typer-0.3.0.tgz",
"integrity": "sha512-dq+qelQ9akHpcOl/gUVRTxVIOkAJ1wR3QAvb4RsVjS8oVoFjDGTc679wJYmUmknUF5HwMLOgb5O+a3KxfWapPQ==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/memoize-one": {
"version": "6.0.0",
"resolved": "https://registry.npmjs.org/memoize-one/-/memoize-one-6.0.0.tgz",
"integrity": "sha512-rkpe71W0N0c0Xz6QD0eJETuWAJGnJ9afsl1srmwPrI+yBCkge5EycXXbYRyvL29zZVUWQCY7InPRCv3GDXuZNw==",
"license": "MIT"
},
"node_modules/merge-descriptors": {
"version": "1.0.3",
"resolved": "https://registry.npmjs.org/merge-descriptors/-/merge-descriptors-1.0.3.tgz",
"integrity": "sha512-gaNvAS7TZ897/rVaZ0nMtAyxNyi/pdbjbAwUpFQpN70GqnVfOiXpeUUMKRBmzXaSQ8DdTX4/0ms62r2K+hE6mQ==",
"dev": true,
"license": "MIT",
"funding": {
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/merge-refs": {
"version": "1.3.0",
"resolved": "https://registry.npmjs.org/merge-refs/-/merge-refs-1.3.0.tgz",
@ -15461,6 +15844,16 @@
"node": ">= 8"
}
},
"node_modules/methods": {
"version": "1.1.2",
"resolved": "https://registry.npmjs.org/methods/-/methods-1.1.2.tgz",
"integrity": "sha512-iclAHeNqNm68zFtnZ0e+1L2yUIdvzNoauKU4WBA3VvH/vPFieF7qfRlwUZU+DA9P9bPXIS90ulxoUoCH23sV2w==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/mhchemparser": {
"version": "4.2.1",
"resolved": "https://registry.npmjs.org/mhchemparser/-/mhchemparser-4.2.1.tgz",
@ -16166,6 +16559,19 @@
"url": "https://github.com/sponsors/aidenybai"
}
},
"node_modules/mime": {
"version": "1.6.0",
"resolved": "https://registry.npmjs.org/mime/-/mime-1.6.0.tgz",
"integrity": "sha512-x0Vn8spI+wuJ1O6S7gnbaQg8Pxh4NNHb7KSINmEWKiPE4RKOplvijn+NkmYmmRgP68mc70j2EbeTFRsrswaQeg==",
"dev": true,
"license": "MIT",
"bin": {
"mime": "cli.js"
},
"engines": {
"node": ">=4"
}
},
"node_modules/mime-db": {
"version": "1.54.0",
"resolved": "https://registry.npmjs.org/mime-db/-/mime-db-1.54.0.tgz",
@ -16636,12 +17042,38 @@
"node": ">= 6"
}
},
"node_modules/object-inspect": {
"version": "1.13.4",
"resolved": "https://registry.npmjs.org/object-inspect/-/object-inspect-1.13.4.tgz",
"integrity": "sha512-W67iLl4J2EXEGTbfeHCffrjDfitvLANg0UlX3wFUUSTx92KXRFegMHUVgSqE+wvhAbi4WqjGg9czysTV2Epbew==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/ohash": {
"version": "1.1.6",
"resolved": "https://registry.npmjs.org/ohash/-/ohash-1.1.6.tgz",
"integrity": "sha512-TBu7PtV8YkAZn0tSxobKY2n2aAQva936lhRrj6957aDaCf9IEtqsKbgMzXE/F/sjqYOwmrukeORHNLe5glk7Cg==",
"license": "MIT"
},
"node_modules/on-finished": {
"version": "2.4.1",
"resolved": "https://registry.npmjs.org/on-finished/-/on-finished-2.4.1.tgz",
"integrity": "sha512-oVlzkg3ENAhCk2zdv7IJwd/QUD4z2RxRwpkcGY8psCVcCYZNq4wYnVWALHM+brtuJjePWiYF/ClmuDr8Ch5+kg==",
"dev": true,
"license": "MIT",
"dependencies": {
"ee-first": "1.1.1"
},
"engines": {
"node": ">= 0.8"
}
},
"node_modules/once": {
"version": "1.4.0",
"resolved": "https://registry.npmjs.org/once/-/once-1.4.0.tgz",
@ -16799,6 +17231,16 @@
"integrity": "sha512-Ofn/CTFzRGTTxwpNEs9PP93gXShHcTq255nzRYSKe8AkVpZY7e1fpmTfOyoIvjP5HG7Z2ZM7VS9PPhQGW2pOpw==",
"license": "MIT"
},
"node_modules/parseurl": {
"version": "1.3.3",
"resolved": "https://registry.npmjs.org/parseurl/-/parseurl-1.3.3.tgz",
"integrity": "sha512-CiyeOxFT/JZyN5m0z9PfXw4SCBJ6Sygz1Dpl0wqjlhDEGGBP1GnsUVEL0p63hoG1fcj3fHynXi9NYO4nWOL+qQ==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/path-exists": {
"version": "4.0.0",
"resolved": "https://registry.npmjs.org/path-exists/-/path-exists-4.0.0.tgz",
@ -16856,6 +17298,13 @@
"integrity": "sha512-JNAzZcXrCt42VGLuYz0zfAzDfAvJWW6AfYlDBQyDV5DClI2m5sAmK+OIO7s59XfsRsWHp02jAJrRadPRGTt6SQ==",
"license": "ISC"
},
"node_modules/path-to-regexp": {
"version": "0.1.12",
"resolved": "https://registry.npmjs.org/path-to-regexp/-/path-to-regexp-0.1.12.tgz",
"integrity": "sha512-RA1GjUVMnvYFxuqovrEqZoxxW5NUZqbwKtYz/Tt7nXerk0LbLblQmrsgdeOxV5SFHf0UDggjS/bSeOZwt1pmEQ==",
"dev": true,
"license": "MIT"
},
"node_modules/path-type": {
"version": "4.0.0",
"resolved": "https://registry.npmjs.org/path-type/-/path-type-4.0.0.tgz",
@ -17357,6 +17806,20 @@
"integrity": "sha512-vtK/94akxsTMhe0/cbfpR+syPuszcuwhqVjJq26CuNDgFGj682oRBXOP5MJpv2r7JtE8MsiepGIqvvOTBwn2vA==",
"license": "ISC"
},
"node_modules/proxy-addr": {
"version": "2.0.7",
"resolved": "https://registry.npmjs.org/proxy-addr/-/proxy-addr-2.0.7.tgz",
"integrity": "sha512-llQsMLSUDUPT44jdrU/O37qlnifitDP+ZwrmmZcoSKyLKvtZxpyV0n2/bD/N4tBAAZ/gJEdZU7KMraoK1+XYAg==",
"dev": true,
"license": "MIT",
"dependencies": {
"forwarded": "0.2.0",
"ipaddr.js": "1.9.1"
},
"engines": {
"node": ">= 0.10"
}
},
"node_modules/proxy-from-env": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/proxy-from-env/-/proxy-from-env-1.1.0.tgz",
@ -17423,6 +17886,22 @@
],
"license": "MIT"
},
"node_modules/qs": {
"version": "6.13.0",
"resolved": "https://registry.npmjs.org/qs/-/qs-6.13.0.tgz",
"integrity": "sha512-+38qI9SOr8tfZ4QmJNplMUxqjbe7LKvvZgWdExBOmd+egZTtjLB67Gu0HRX3u/XOq7UU2Nx6nsjvS16Z9uwfpg==",
"dev": true,
"license": "BSD-3-Clause",
"dependencies": {
"side-channel": "^1.0.6"
},
"engines": {
"node": ">=0.6"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/querystringify": {
"version": "2.2.0",
"resolved": "https://registry.npmjs.org/querystringify/-/querystringify-2.2.0.tgz",
@ -17462,6 +17941,49 @@
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/range-parser": {
"version": "1.2.1",
"resolved": "https://registry.npmjs.org/range-parser/-/range-parser-1.2.1.tgz",
"integrity": "sha512-Hrgsx+orqoygnmhFbKaHE6c296J+HTAQXoxEF6gNupROmmGJRoyzfG3ccAveqCBrwr/2yxQ5BVd/GTl5agOwSg==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.6"
}
},
"node_modules/raw-body": {
"version": "3.0.1",
"resolved": "https://registry.npmjs.org/raw-body/-/raw-body-3.0.1.tgz",
"integrity": "sha512-9G8cA+tuMS75+6G/TzW8OtLzmBDMo8p1JRxN5AZ+LAp8uxGA8V8GZm4GQ4/N5QNQEnLmg6SS7wyuSmbKepiKqA==",
"dev": true,
"license": "MIT",
"dependencies": {
"bytes": "3.1.2",
"http-errors": "2.0.0",
"iconv-lite": "0.7.0",
"unpipe": "1.0.0"
},
"engines": {
"node": ">= 0.10"
}
},
"node_modules/raw-body/node_modules/iconv-lite": {
"version": "0.7.0",
"resolved": "https://registry.npmjs.org/iconv-lite/-/iconv-lite-0.7.0.tgz",
"integrity": "sha512-cf6L2Ds3h57VVmkZe+Pn+5APsT7FpqJtEhhieDCvrE2MK5Qk9MyffgQyuxQTm6BChfeZNtcOLHp9IcWRVcIcBQ==",
"dev": true,
"license": "MIT",
"dependencies": {
"safer-buffer": ">= 2.1.2 < 3.0.0"
},
"engines": {
"node": ">=0.10.0"
},
"funding": {
"type": "opencollective",
"url": "https://opencollective.com/express"
}
},
"node_modules/rc": {
"version": "1.2.8",
"resolved": "https://registry.npmjs.org/rc/-/rc-1.2.8.tgz",
@ -18797,6 +19319,74 @@
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/send": {
"version": "0.19.0",
"resolved": "https://registry.npmjs.org/send/-/send-0.19.0.tgz",
"integrity": "sha512-dW41u5VfLXu8SJh5bwRmyYUbAoSB3c9uQh6L8h/KtsFREPWpbX1lrljJo186Jc4nmci/sGUZ9a0a0J2zgfq2hw==",
"dev": true,
"license": "MIT",
"dependencies": {
"debug": "2.6.9",
"depd": "2.0.0",
"destroy": "1.2.0",
"encodeurl": "~1.0.2",
"escape-html": "~1.0.3",
"etag": "~1.8.1",
"fresh": "0.5.2",
"http-errors": "2.0.0",
"mime": "1.6.0",
"ms": "2.1.3",
"on-finished": "2.4.1",
"range-parser": "~1.2.1",
"statuses": "2.0.1"
},
"engines": {
"node": ">= 0.8.0"
}
},
"node_modules/send/node_modules/debug": {
"version": "2.6.9",
"resolved": "https://registry.npmjs.org/debug/-/debug-2.6.9.tgz",
"integrity": "sha512-bC7ElrdJaJnPbAP+1EotYvqZsb3ecl5wi6Bfi6BJTUcNowp6cvspg0jXznRTKDjm/E7AdgFBVeAPVMNcKGsHMA==",
"dev": true,
"license": "MIT",
"dependencies": {
"ms": "2.0.0"
}
},
"node_modules/send/node_modules/debug/node_modules/ms": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/ms/-/ms-2.0.0.tgz",
"integrity": "sha512-Tpp60P6IUJDTuOq/5Z8cdskzJujfwqfOTkrwIwj7IRISpnkJnT6SyJ4PCPnGMoFjC9ddhal5KVIYtAt97ix05A==",
"dev": true,
"license": "MIT"
},
"node_modules/send/node_modules/encodeurl": {
"version": "1.0.2",
"resolved": "https://registry.npmjs.org/encodeurl/-/encodeurl-1.0.2.tgz",
"integrity": "sha512-TPJXq8JqFaVYm2CWmPvnP2Iyo4ZSM7/QKcSmuMLDObfpH5fi7RUGmd/rTDf+rut/saiDiQEeVTNgAmJEdAOx0w==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/serve-static": {
"version": "1.16.2",
"resolved": "https://registry.npmjs.org/serve-static/-/serve-static-1.16.2.tgz",
"integrity": "sha512-VqpjJZKadQB/PEbEwvFdO43Ax5dFBZ2UECszz8bQ7pi7wt//PWe1P6MN7eCnjsatYtBT6EuiClbjSWP2WrIoTw==",
"dev": true,
"license": "MIT",
"dependencies": {
"encodeurl": "~2.0.0",
"escape-html": "~1.0.3",
"parseurl": "~1.3.3",
"send": "0.19.0"
},
"engines": {
"node": ">= 0.8.0"
}
},
"node_modules/set-blocking": {
"version": "2.0.0",
"resolved": "https://registry.npmjs.org/set-blocking/-/set-blocking-2.0.0.tgz",
@ -18805,6 +19395,13 @@
"optional": true,
"peer": true
},
"node_modules/setprototypeof": {
"version": "1.2.0",
"resolved": "https://registry.npmjs.org/setprototypeof/-/setprototypeof-1.2.0.tgz",
"integrity": "sha512-E5LDX7Wrp85Kil5bhZv46j8jOeboKq5JMmYM3gVGdGH8xFpPWXUMsNrlODCrkoxMEeNi/XZIwuRvY4XNwYMJpw==",
"dev": true,
"license": "ISC"
},
"node_modules/shadcn-ui": {
"version": "0.9.5",
"resolved": "https://registry.npmjs.org/shadcn-ui/-/shadcn-ui-0.9.5.tgz",
@ -18860,6 +19457,82 @@
"suid": "bin/short-unique-id"
}
},
"node_modules/side-channel": {
"version": "1.1.0",
"resolved": "https://registry.npmjs.org/side-channel/-/side-channel-1.1.0.tgz",
"integrity": "sha512-ZX99e6tRweoUXqR+VBrslhda51Nh5MTQwou5tnUDgbtyM0dBgmhEDtWGP/xbKn6hqfPRHujUNwz5fy/wbbhnpw==",
"dev": true,
"license": "MIT",
"dependencies": {
"es-errors": "^1.3.0",
"object-inspect": "^1.13.3",
"side-channel-list": "^1.0.0",
"side-channel-map": "^1.0.1",
"side-channel-weakmap": "^1.0.2"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/side-channel-list": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/side-channel-list/-/side-channel-list-1.0.0.tgz",
"integrity": "sha512-FCLHtRD/gnpCiCHEiJLOwdmFP+wzCmDEkc9y7NsYxeF4u7Btsn1ZuwgwJGxImImHicJArLP4R0yX4c2KCrMrTA==",
"dev": true,
"license": "MIT",
"dependencies": {
"es-errors": "^1.3.0",
"object-inspect": "^1.13.3"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/side-channel-map": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/side-channel-map/-/side-channel-map-1.0.1.tgz",
"integrity": "sha512-VCjCNfgMsby3tTdo02nbjtM/ewra6jPHmpThenkTYh8pG9ucZ/1P8So4u4FGBek/BjpOVsDCMoLA/iuBKIFXRA==",
"dev": true,
"license": "MIT",
"dependencies": {
"call-bound": "^1.0.2",
"es-errors": "^1.3.0",
"get-intrinsic": "^1.2.5",
"object-inspect": "^1.13.3"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/side-channel-weakmap": {
"version": "1.0.2",
"resolved": "https://registry.npmjs.org/side-channel-weakmap/-/side-channel-weakmap-1.0.2.tgz",
"integrity": "sha512-WPS/HvHQTYnHisLo9McqBHOJk2FkHO/tlpvldyrnem4aeQp4hai3gythswg6p01oSoTl58rcpiFAjF2br2Ak2A==",
"dev": true,
"license": "MIT",
"dependencies": {
"call-bound": "^1.0.2",
"es-errors": "^1.3.0",
"get-intrinsic": "^1.2.5",
"object-inspect": "^1.13.3",
"side-channel-map": "^1.0.1"
},
"engines": {
"node": ">= 0.4"
},
"funding": {
"url": "https://github.com/sponsors/ljharb"
}
},
"node_modules/signal-exit": {
"version": "3.0.7",
"resolved": "https://registry.npmjs.org/signal-exit/-/signal-exit-3.0.7.tgz",
@ -19195,6 +19868,16 @@
"node": ">=8"
}
},
"node_modules/statuses": {
"version": "2.0.1",
"resolved": "https://registry.npmjs.org/statuses/-/statuses-2.0.1.tgz",
"integrity": "sha512-RwNA9Z/7PrK06rYLIzFMlaF+l73iwpzsqRIFgbMLbTcLD6cOao82TaWefPXQvB2fOC4AjuYSEndS7N/mTCbkdQ==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/streamx": {
"version": "2.22.1",
"resolved": "https://registry.npmjs.org/streamx/-/streamx-2.22.1.tgz",
@ -19886,6 +20569,16 @@
"node": ">=8.0"
}
},
"node_modules/toidentifier": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/toidentifier/-/toidentifier-1.0.1.tgz",
"integrity": "sha512-o5sSPKEkg/DIQNmH43V0/uerLrpzVedkUh8tGNvaeXpfpuwjKenlSox/2O/BTlZUtEe+JG7s5YhEz608PlAHRA==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">=0.6"
}
},
"node_modules/token-types": {
"version": "6.0.0",
"resolved": "https://registry.npmjs.org/token-types/-/token-types-6.0.0.tgz",
@ -20075,6 +20768,20 @@
"url": "https://github.com/sponsors/sindresorhus"
}
},
"node_modules/type-is": {
"version": "1.6.18",
"resolved": "https://registry.npmjs.org/type-is/-/type-is-1.6.18.tgz",
"integrity": "sha512-TkRKr9sUTxEH8MdfuCSP7VizJyzRNMjj2J2do2Jr3Kym598JVdEksuzPQCnlFPW4ky9Q+iA+ma9BGm06XQBy8g==",
"dev": true,
"license": "MIT",
"dependencies": {
"media-typer": "0.3.0",
"mime-types": "~2.1.24"
},
"engines": {
"node": ">= 0.6"
}
},
"node_modules/typedarray-to-buffer": {
"version": "3.1.5",
"resolved": "https://registry.npmjs.org/typedarray-to-buffer/-/typedarray-to-buffer-3.1.5.tgz",
@ -20380,6 +21087,16 @@
"node": ">= 4.0.0"
}
},
"node_modules/unpipe": {
"version": "1.0.0",
"resolved": "https://registry.npmjs.org/unpipe/-/unpipe-1.0.0.tgz",
"integrity": "sha512-pjy2bYhSsufwWlKwPc+l3cN7+wuJlK6uz0YdJEOlQDbl6jo/YlPi4mb8agUkVC8BF7V8NuzeyPNqRksA3hztKQ==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.8"
}
},
"node_modules/unplugin": {
"version": "1.16.1",
"resolved": "https://registry.npmjs.org/unplugin/-/unplugin-1.16.1.tgz",
@ -20561,6 +21278,16 @@
"integrity": "sha512-EPD5q1uXyFxJpCrLnCc1nHnq3gOa6DZBocAIiI2TaSCA7VCJ1UJDMagCzIkXNsUYfD1daK//LTEQ8xiIbrHtcw==",
"license": "MIT"
},
"node_modules/utils-merge": {
"version": "1.0.1",
"resolved": "https://registry.npmjs.org/utils-merge/-/utils-merge-1.0.1.tgz",
"integrity": "sha512-pMZTvIkT1d+TFGvDOqodOclx0QWkkgi6Tdoa8gC8ffGAAqz9pzPTZWAybbsHHoED/ztMtkv/VoYTYyShUn81hA==",
"dev": true,
"license": "MIT",
"engines": {
"node": ">= 0.4.0"
}
},
"node_modules/uuid": {
"version": "10.0.0",
"resolved": "https://registry.npmjs.org/uuid/-/uuid-10.0.0.tgz",
@ -21759,6 +22486,16 @@
"url": "https://github.com/sponsors/colinhacks"
}
},
"node_modules/zod-to-json-schema": {
"version": "3.24.6",
"resolved": "https://registry.npmjs.org/zod-to-json-schema/-/zod-to-json-schema-3.24.6.tgz",
"integrity": "sha512-h/z3PKvcTcTetyjl1fkj79MHNEjm+HpD6NXheWjzOekY7kV+lwDYnHw+ivHkijnCSMz1yJaWBD9vu/Fcmk+vEg==",
"dev": true,
"license": "ISC",
"peerDependencies": {
"zod": "^3.24.1"
}
},
"node_modules/zustand": {
"version": "4.5.7",
"resolved": "https://registry.npmjs.org/zustand/-/zustand-4.5.7.tgz",

View File

@ -118,6 +118,7 @@
"devDependencies": {
"@biomejs/biome": "2.1.1",
"@jest/types": "^30.0.1",
"@modelcontextprotocol/server-everything": "^0.6.2",
"@playwright/test": "^1.52.0",
"@swc/cli": "^0.5.2",
"@swc/core": "^1.6.1",

View File

@ -49,7 +49,7 @@ export default function AddMcpServerModal({
: useState(false);
const [type, setType] = useState(
initialData ? (initialData.command ? "STDIO" : "SSE") : "JSON",
initialData ? (initialData.command ? "STDIO" : "HTTP") : "JSON",
);
const [jsonValue, setJsonValue] = useState("");
const [error, setError] = useState<string | null>(
@ -73,10 +73,10 @@ export default function AddMcpServerModal({
setStdioCommand("");
setStdioArgs([""]);
setStdioEnv([]);
setSseName("");
setSseUrl("");
setSseEnv([]);
setSseHeaders([]);
setHttpName("");
setHttpUrl("");
setHttpEnv([]);
setHttpHeaders([]);
};
// STDIO state
@ -87,25 +87,27 @@ export default function AddMcpServerModal({
);
const [stdioEnv, setStdioEnv] = useState<any>(initialData?.env || []);
// SSE state
const [sseName, setSseName] = useState(initialData?.name || "");
const [sseUrl, setSseUrl] = useState(initialData?.url || "");
const [sseEnv, setSseEnv] = useState<any>(initialData?.env || []);
const [sseHeaders, setSseHeaders] = useState<any>(initialData?.headers || []);
// HTTP state
const [httpName, setHttpName] = useState(initialData?.name || "");
const [httpUrl, setHttpUrl] = useState(initialData?.url || "");
const [httpEnv, setHttpEnv] = useState<any>(initialData?.env || []);
const [httpHeaders, setHttpHeaders] = useState<any>(
initialData?.headers || [],
);
useEffect(() => {
if (open) {
setType(initialData ? (initialData.command ? "STDIO" : "SSE") : "JSON");
setType(initialData ? (initialData.command ? "STDIO" : "HTTP") : "JSON");
setError(null);
setJsonValue("");
setStdioName(initialData?.name || "");
setStdioCommand(initialData?.command || "");
setStdioArgs(initialData?.args || [""]);
setStdioEnv(initialData?.env || []);
setSseName(initialData?.name || "");
setSseUrl(initialData?.url || "");
setSseEnv(initialData?.env || []);
setSseHeaders(initialData?.headers || []);
setHttpName(initialData?.name || "");
setHttpUrl(initialData?.url || "");
setHttpEnv(initialData?.env || []);
setHttpHeaders(initialData?.headers || []);
}
}, [open]);
@ -159,12 +161,12 @@ export default function AddMcpServerModal({
}
return;
}
if (type === "SSE") {
if (!sseName.trim() || !sseUrl.trim()) {
if (type === "HTTP") {
if (!httpName.trim() || !httpUrl.trim()) {
setError("Name and URL are required.");
return;
}
const name = parseString(sseName, [
const name = parseString(httpName, [
"mcp_name_case",
"no_blank",
"lowercase",
@ -172,9 +174,9 @@ export default function AddMcpServerModal({
try {
await modifyMCPServer({
name,
env: parseEnvList(sseEnv),
url: sseUrl,
headers: parseEnvList(sseHeaders),
env: parseEnvList(httpEnv),
url: httpUrl,
headers: parseEnvList(httpHeaders),
});
if (!initialData) {
await queryClient.setQueryData(["useGetMCPServers"], (old: any) => {
@ -183,10 +185,10 @@ export default function AddMcpServerModal({
}
onSuccess?.(name);
setOpen(false);
setSseName("");
setSseUrl("");
setSseEnv([]);
setSseHeaders([]);
setHttpName("");
setHttpUrl("");
setHttpEnv([]);
setHttpHeaders([]);
setError(null);
} catch (err: any) {
setError(err?.message || "Failed to add MCP server.");
@ -285,11 +287,11 @@ export default function AddMcpServerModal({
STDIO
</TabsTrigger>
<TabsTrigger
data-testid="sse-tab"
disabled={!!initialData && type !== "SSE"}
value="SSE"
data-testid="http-tab"
disabled={!!initialData && type !== "HTTP"}
value="HTTP"
>
SSE
Streamable HTTP/SSE
</TabsTrigger>
</TabsList>
</div>
@ -374,52 +376,53 @@ export default function AddMcpServerModal({
</div>
</div>
</TabsContent>
<TabsContent value="SSE">
<TabsContent value="HTTP">
<div className="flex flex-col gap-4">
<div className="flex flex-col gap-2">
<Label className="flex items-start gap-1 !text-mmd">
Name<span className="text-red-500">*</span>
</Label>
<Input
value={sseName}
onChange={(e) => setSseName(e.target.value)}
value={httpName}
onChange={(e) => setHttpName(e.target.value)}
placeholder="Name"
data-testid="sse-name-input"
data-testid="http-name-input"
disabled={isPending}
/>
</div>
<div className="flex flex-col gap-2">
<Label className="flex items-start gap-1 !text-mmd">
SSE URL<span className="text-red-500">*</span>
Streamable HTTP/SSE URL
<span className="text-red-500">*</span>
</Label>
<Input
value={sseUrl}
onChange={(e) => setSseUrl(e.target.value)}
placeholder="SSE URL"
data-testid="sse-url-input"
value={httpUrl}
onChange={(e) => setHttpUrl(e.target.value)}
placeholder="Streamable HTTP/SSE URL"
data-testid="http-url-input"
disabled={isPending}
/>
</div>
<div className="flex flex-col gap-2">
<Label className="!text-mmd">Headers</Label>
<IOKeyPairInput
value={sseHeaders}
onChange={setSseHeaders}
value={httpHeaders}
onChange={setHttpHeaders}
duplicateKey={false}
isList={true}
isInputField={true}
testId="sse-headers"
testId="http-headers"
/>
</div>
<div className="flex flex-col gap-2">
<Label className="!text-mmd">Environment Variables</Label>
<IOKeyPairInput
value={sseEnv}
onChange={setSseEnv}
value={httpEnv}
onChange={setHttpEnv}
duplicateKey={false}
isList={true}
isInputField={true}
testId="sse-env"
testId="http-env"
/>
</div>
</div>

View File

@ -159,7 +159,7 @@ test(
timeout: 3000,
});
await expect(page.getByTestId("sse-tab")).toBeDisabled({
await expect(page.getByTestId("http-tab")).toBeDisabled({
timeout: 3000,
});
@ -242,7 +242,7 @@ test(
const sidebarButton = page.getByTestId("sidebar-add-mcp-server-button");
const fallbackButton = page.getByTestId("add-mcp-server-button-sidebar");
if (await sidebarButton.isVisible({ timeout: 3000 }).catch(() => false)) {
if (await sidebarButton.isVisible({ timeout: 5000 }).catch(() => false)) {
await sidebarButton.click();
} else {
await fallbackButton.click();
@ -267,6 +267,8 @@ test(
await page.getByTestId("add-mcp-server-button").click();
await page.waitForTimeout(1000);
await page.getByTestId(`add-component-button-${testName}`).click();
await expect(page.getByTestId("dropdown_str_tool")).toBeVisible({
@ -521,7 +523,7 @@ test(
);
test(
"SSE MCP server fields should persist after saving and editing",
"HTTP/SSE MCP server fields should persist after saving and editing",
{ tag: ["@release", "@workspace", "@components"] },
async ({ page }) => {
await awaitBootstrapTest(page);
@ -560,16 +562,16 @@ test(
timeout: 30000,
});
// Go to SSE tab and fill all fields
await page.getByTestId("sse-tab").click();
await page.waitForSelector('[data-testid="sse-name-input"]', {
// Go to HTTP tab and fill all fields
await page.getByTestId("http-tab").click();
await page.waitForSelector('[data-testid="http-name-input"]', {
state: "visible",
timeout: 30000,
});
// Test data with random suffix
const randomSuffix = Math.floor(Math.random() * 90000) + 10000; // 5-digit random number
const testName = `test_sse_server_${randomSuffix}`;
const testName = `test_http_server_${randomSuffix}`;
const testUrl = "https://api.example.com/mcp";
const testHeaderKey1 = "Authorization";
const testHeaderValue1 = "Bearer token123";
@ -581,26 +583,26 @@ test(
const testEnvValue2 = "3";
// Fill basic fields
await page.getByTestId("sse-name-input").fill(testName);
await page.getByTestId("sse-url-input").fill(testUrl);
await page.getByTestId("http-name-input").fill(testName);
await page.getByTestId("http-url-input").fill(testUrl);
// Add first header
await page.getByTestId("sse-headers-key-0").fill(testHeaderKey1);
await page.getByTestId("sse-headers-value-0").fill(testHeaderValue1);
await page.getByTestId("http-headers-key-0").fill(testHeaderKey1);
await page.getByTestId("http-headers-value-0").fill(testHeaderValue1);
// Add second header
await page.getByTestId("sse-headers-plus-btn-0").click();
await page.getByTestId("sse-headers-key-1").fill(testHeaderKey2);
await page.getByTestId("sse-headers-value-1").fill(testHeaderValue2);
await page.getByTestId("http-headers-plus-btn-0").click();
await page.getByTestId("http-headers-key-1").fill(testHeaderKey2);
await page.getByTestId("http-headers-value-1").fill(testHeaderValue2);
// Add first environment variable
await page.getByTestId("sse-env-key-0").fill(testEnvKey1);
await page.getByTestId("sse-env-value-0").fill(testEnvValue1);
await page.getByTestId("http-env-key-0").fill(testEnvKey1);
await page.getByTestId("http-env-value-0").fill(testEnvValue1);
// Add second environment variable
await page.getByTestId("sse-env-plus-btn-0").click();
await page.getByTestId("sse-env-key-1").fill(testEnvKey2);
await page.getByTestId("sse-env-value-1").fill(testEnvValue2);
await page.getByTestId("http-env-plus-btn-0").click();
await page.getByTestId("http-env-key-1").fill(testEnvKey2);
await page.getByTestId("http-env-value-1").fill(testEnvValue2);
// Save the server
await page.getByTestId("add-mcp-server-button").click();
@ -641,32 +643,32 @@ test(
});
// Verify all fields persisted correctly
expect(await page.getByTestId("sse-name-input").inputValue()).toBe(
expect(await page.getByTestId("http-name-input").inputValue()).toBe(
testName,
);
expect(await page.getByTestId("sse-url-input").inputValue()).toBe(testUrl);
expect(await page.getByTestId("sse-headers-key-0").inputValue()).toBe(
expect(await page.getByTestId("http-url-input").inputValue()).toBe(testUrl);
expect(await page.getByTestId("http-headers-key-0").inputValue()).toBe(
testHeaderKey1,
);
expect(await page.getByTestId("sse-headers-value-0").inputValue()).toBe(
expect(await page.getByTestId("http-headers-value-0").inputValue()).toBe(
testHeaderValue1,
);
expect(await page.getByTestId("sse-headers-key-1").inputValue()).toBe(
expect(await page.getByTestId("http-headers-key-1").inputValue()).toBe(
testHeaderKey2,
);
expect(await page.getByTestId("sse-headers-value-1").inputValue()).toBe(
expect(await page.getByTestId("http-headers-value-1").inputValue()).toBe(
testHeaderValue2,
);
expect(await page.getByTestId("sse-env-key-0").inputValue()).toBe(
expect(await page.getByTestId("http-env-key-0").inputValue()).toBe(
testEnvKey1,
);
expect(await page.getByTestId("sse-env-value-0").inputValue()).toBe(
expect(await page.getByTestId("http-env-value-0").inputValue()).toBe(
testEnvValue1,
);
expect(await page.getByTestId("sse-env-key-1").inputValue()).toBe(
expect(await page.getByTestId("http-env-key-1").inputValue()).toBe(
testEnvKey2,
);
expect(await page.getByTestId("sse-env-value-1").inputValue()).toBe(
expect(await page.getByTestId("http-env-value-1").inputValue()).toBe(
testEnvValue2,
);
@ -838,7 +840,7 @@ test(
timeout: 3000,
});
await expect(page.getByTestId("sse-tab")).toBeDisabled({
await expect(page.getByTestId("http-tab")).toBeDisabled({
timeout: 3000,
});
@ -982,3 +984,214 @@ test(
expect(fetchOptionCount2).toBeGreaterThan(0);
},
);
test(
"Streamable HTTP MCP server with server-everything should load tools correctly",
{ tag: ["@release", "@workspace", "@components"] },
async ({ page }) => {
// Start the MCP server with proper health checking
const server = "https://mcp.deepwiki.com/mcp";
await awaitBootstrapTest(page);
await page.waitForSelector('[data-testid="blank-flow"]', {
timeout: 30000,
});
await page.getByTestId("blank-flow").click();
await page.getByTestId("sidebar-search-input").click();
await page.getByTestId("sidebar-search-input").fill("mcp tools");
await page.waitForSelector('[data-testid="agentsMCP Tools"]', {
timeout: 30000,
});
await page
.getByTestId("agentsMCP Tools")
.dragTo(page.locator('//*[@id="react-flow-id"]'), {
targetPosition: { x: 100, y: 100 },
});
await adjustScreenView(page, { numberOfZoomOut: 3 });
try {
await page.getByText("Add MCP Server", { exact: true }).click({
timeout: 5000,
});
} catch (_error) {
await page.getByTestId("mcp-server-dropdown").click({ timeout: 3000 });
await page.getByText("Add MCP Server", { exact: true }).click({
timeout: 5000,
});
}
await page.waitForSelector('[data-testid="add-mcp-server-button"]', {
state: "visible",
timeout: 30000,
});
// Switch to HTTP tab for Streamable HTTP
await page.getByTestId("http-tab").click();
await page.waitForSelector('[data-testid="http-name-input"]', {
state: "visible",
timeout: 30000,
});
const randomSuffix = Math.floor(Math.random() * 90000) + 10000;
const testName = `test_streamable_http_${randomSuffix}`;
// Fill in the server details
await page.getByTestId("http-name-input").fill(testName);
// Use the HTTP endpoint URL
await page.getByTestId("http-url-input").fill(server);
await page.getByTestId("add-mcp-server-button").click();
// Wait for tools to load with proper timeout
await page.waitForSelector(
'[data-testid="dropdown_str_tool"]:not([disabled])',
{
timeout: 10000,
state: "visible",
},
);
await page.getByTestId("dropdown_str_tool").click();
// Check for tools from server
const toolOptions = page.locator('[data-testid*="-option"]');
const toolCount = await toolOptions.count();
// server-everything should have multiple tools (at least 5+)
expect(toolCount).toBeGreaterThan(5);
// Verify specific tools exist from server-everything
const readWikiStructureOption = page.getByTestId(
"read_wiki_structure-0-option",
);
expect(await readWikiStructureOption.count()).toBeGreaterThan(0);
// Select the option to verify it loads properly
await readWikiStructureOption.last().click();
// Wait for the tool input field to appear
await page.waitForSelector(
'[data-testid="popover-anchor-input-repoName"]',
{
state: "visible",
timeout: 10000,
},
);
// Verify the input field is present
await expect(
page.getByTestId("popover-anchor-input-repoName"),
).toBeVisible();
},
);
test(
"SSE MCP server with deepwiki should load tools correctly",
{ tag: ["@release", "@workspace", "@components"] },
async ({ page }) => {
// Start the MCP server with proper health checking
const server = "https://mcp.deepwiki.com/sse";
await awaitBootstrapTest(page);
await page.waitForSelector('[data-testid="blank-flow"]', {
timeout: 30000,
});
await page.getByTestId("blank-flow").click();
await page.getByTestId("sidebar-search-input").click();
await page.getByTestId("sidebar-search-input").fill("mcp tools");
await page.waitForSelector('[data-testid="agentsMCP Tools"]', {
timeout: 30000,
});
await page
.getByTestId("agentsMCP Tools")
.dragTo(page.locator('//*[@id="react-flow-id"]'), {
targetPosition: { x: 100, y: 100 },
});
await adjustScreenView(page, { numberOfZoomOut: 3 });
try {
await page.getByText("Add MCP Server", { exact: true }).click({
timeout: 5000,
});
} catch (_error) {
await page.getByTestId("mcp-server-dropdown").click({ timeout: 3000 });
await page.getByText("Add MCP Server", { exact: true }).click({
timeout: 5000,
});
}
await page.waitForSelector('[data-testid="add-mcp-server-button"]', {
state: "visible",
timeout: 30000,
});
// Switch to HTTP tab for SSE
await page.getByTestId("http-tab").click();
await page.waitForSelector('[data-testid="http-name-input"]', {
state: "visible",
timeout: 30000,
});
const randomSuffix = Math.floor(Math.random() * 90000) + 10000;
const testName = `test_sse_${randomSuffix}`;
// Fill in the server details
await page.getByTestId("http-name-input").fill(testName);
// Use the HTTP endpoint URL
await page.getByTestId("http-url-input").fill(server);
await page.getByTestId("add-mcp-server-button").click();
// Wait for tools to load with proper timeout
await page.waitForSelector(
'[data-testid="dropdown_str_tool"]:not([disabled])',
{
timeout: 10000,
state: "visible",
},
);
await page.getByTestId("dropdown_str_tool").click();
// Check for tools from wiki
const toolOptions = page.locator('[data-testid*="-option"]');
const toolCount = await toolOptions.count();
// server-everything should have multiple tools (at least 5+)
expect(toolCount).toBeGreaterThan(5);
// Verify specific tools exist from server-everything
const readWikiStructureOption = page.getByTestId(
"read_wiki_structure-0-option",
);
expect(await readWikiStructureOption.count()).toBeGreaterThan(0);
// Select the readWikiStructure to verify it loads properly
await readWikiStructureOption.last().click();
// Wait for the tool input field to appear
await page.waitForSelector(
'[data-testid="popover-anchor-input-repoName"]',
{
state: "visible",
timeout: 10000,
},
);
// Verify the input field is present
await expect(
page.getByTestId("popover-anchor-input-repoName"),
).toBeVisible();
},
);

View File

@ -80,7 +80,7 @@ async def handle_on_chain_start(
header={"title": "Input", "icon": "MessageSquare"},
)
agent_message.content_blocks[0].contents.append(text_content)
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
start_time = perf_counter()
return agent_message, start_time
@ -151,7 +151,7 @@ async def handle_on_chain_end(
header={"title": "Output", "icon": "MessageSquare"},
)
agent_message.content_blocks[0].contents.append(text_content)
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
start_time = perf_counter()
return agent_message, start_time
@ -190,7 +190,7 @@ async def handle_on_tool_start(
tool_blocks_map[tool_key] = tool_content
agent_message.content_blocks[0].contents.append(tool_content)
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
if agent_message.content_blocks and agent_message.content_blocks[0].contents:
tool_blocks_map[tool_key] = agent_message.content_blocks[0].contents[-1]
return agent_message, new_start_time
@ -210,7 +210,7 @@ async def handle_on_tool_end(
if tool_content and isinstance(tool_content, ToolContent):
# Call send_message_method first to get the updated message structure
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
new_start_time = perf_counter()
# Now find and update the tool content in the current message
@ -258,7 +258,7 @@ async def handle_on_tool_error(
tool_content.error = event["data"].get("error", "Unknown error")
tool_content.duration = _calculate_duration(start_time)
tool_content.header = {"title": f"Error using **{tool_content.name}**", "icon": "Hammer"}
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
start_time = perf_counter()
return agent_message, start_time
@ -275,14 +275,14 @@ async def handle_on_chain_stream(
if output and isinstance(output, str | list):
agent_message.text = _extract_output_text(output)
agent_message.properties.state = "complete"
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
start_time = perf_counter()
elif isinstance(data_chunk, AIMessageChunk):
output_text = _extract_output_text(data_chunk.content)
if output_text and isinstance(agent_message.text, str):
agent_message.text += output_text
agent_message.properties.state = "partial"
agent_message = await send_message_method(message=agent_message)
agent_message = await send_message_method(message=agent_message, skip_db_update=True)
if not agent_message.text:
start_time = perf_counter()
return agent_message, start_time
@ -346,13 +346,17 @@ async def process_agent_events(
async for event in agent_executor:
if event["event"] in TOOL_EVENT_HANDLERS:
tool_handler = TOOL_EVENT_HANDLERS[event["event"]]
# Use skip_db_update=True during streaming to avoid DB round-trips
agent_message, start_time = await tool_handler(
event, agent_message, tool_blocks_map, send_message_method, start_time
)
elif event["event"] in CHAIN_EVENT_HANDLERS:
chain_handler = CHAIN_EVENT_HANDLERS[event["event"]]
# Use skip_db_update=True during streaming to avoid DB round-trips
agent_message, start_time = await chain_handler(event, agent_message, send_message_method, start_time)
agent_message.properties.state = "complete"
# Final DB update with the complete message (skip_db_update=False by default)
agent_message = await send_message_method(message=agent_message)
except Exception as e:
raise ExceptionWithMessageError(agent_message, str(e)) from e
return await Message.create(**agent_message.model_dump())

View File

@ -28,8 +28,12 @@ HTTP_ERROR_STATUS_CODE = httpx_codes.BAD_REQUEST # HTTP status code for client
# HTTP status codes used in validation
HTTP_NOT_FOUND = 404
HTTP_METHOD_NOT_ALLOWED = 405
HTTP_NOT_ACCEPTABLE = 406
HTTP_BAD_REQUEST = 400
HTTP_INTERNAL_SERVER_ERROR = 500
HTTP_UNAUTHORIZED = 401
HTTP_FORBIDDEN = 403
# MCP Session Manager constants
settings = get_settings_service().settings
@ -378,8 +382,8 @@ def _validate_node_installation(command: str) -> str:
async def _validate_connection_params(mode: str, command: str | None = None, url: str | None = None) -> None:
"""Validate connection parameters based on mode."""
if mode not in ["Stdio", "SSE"]:
msg = f"Invalid mode: {mode}. Must be either 'Stdio' or 'SSE'"
if mode not in ["Stdio", "Streamable_HTTP", "SSE"]:
msg = f"Invalid mode: {mode}. Must be either 'Stdio', 'Streamable_HTTP', or 'SSE'"
raise ValueError(msg)
if mode == "Stdio" and not command:
@ -387,8 +391,8 @@ async def _validate_connection_params(mode: str, command: str | None = None, url
raise ValueError(msg)
if mode == "Stdio" and command:
_validate_node_installation(command)
if mode == "SSE" and not url:
msg = "URL is required for SSE mode"
if mode in ["Streamable_HTTP", "SSE"] and not url:
msg = f"URL is required for {mode} mode"
raise ValueError(msg)
@ -400,6 +404,7 @@ class MCPSessionManager:
2. Maximum session limits per server to prevent resource exhaustion
3. Idle timeout for automatic session cleanup
4. Periodic cleanup of stale sessions
5. Transport preference caching to avoid retrying failed transports
"""
def __init__(self):
@ -410,6 +415,9 @@ class MCPSessionManager:
self._context_to_session: dict[str, tuple[str, str]] = {}
# Reference count for each active (server_key, session_id)
self._session_refcount: dict[tuple[str, str], int] = {}
# Cache which transport works for each server to avoid retrying failed transports
# server_key -> "streamable_http" | "sse"
self._transport_preference: dict[str, str] = {}
self._cleanup_task = None
self._start_cleanup_task()
@ -467,15 +475,16 @@ class MCPSessionManager:
env_str = str(sorted((connection_params.env or {}).items()))
key_input = f"{command_str}|{env_str}"
return f"stdio_{hash(key_input)}"
elif transport_type == "sse" and (isinstance(connection_params, dict) and "url" in connection_params):
elif transport_type == "streamable_http" and (
isinstance(connection_params, dict) and "url" in connection_params
):
# Include URL and headers for uniqueness
url = connection_params["url"]
headers = str(sorted((connection_params.get("headers", {})).items()))
key_input = f"{url}|{headers}"
return f"sse_{hash(key_input)}"
return f"streamable_http_{hash(key_input)}"
# Fallback to a generic key
# TODO: add option for streamable HTTP in future.
return f"{transport_type}_{hash(str(connection_params))}"
async def _validate_session_connectivity(self, session) -> bool:
@ -525,7 +534,7 @@ class MCPSessionManager:
"""Get or create a session with improved reuse strategy.
The key insight is that we should reuse sessions based on the server
identity (command + args for stdio, URL for SSE) rather than the context_id.
identity (command + args for stdio, URL for Streamable HTTP) rather than the context_id.
This prevents creating a new subprocess for each unique context.
"""
server_key = self._get_server_key(connection_params, transport_type)
@ -578,17 +587,24 @@ class MCPSessionManager:
if transport_type == "stdio":
session, task = await self._create_stdio_session(session_id, connection_params)
elif transport_type == "sse":
session, task = await self._create_sse_session(session_id, connection_params)
actual_transport = "stdio"
elif transport_type == "streamable_http":
# Pass the cached transport preference if available
preferred_transport = self._transport_preference.get(server_key)
session, task, actual_transport = await self._create_streamable_http_session(
session_id, connection_params, preferred_transport
)
# Cache the transport that worked for future connections
self._transport_preference[server_key] = actual_transport
else:
msg = f"Unknown transport type: {transport_type}"
raise ValueError(msg)
# Store session info
# Store session info with the actual transport used
sessions[session_id] = {
"session": session,
"task": task,
"type": transport_type,
"type": actual_transport,
"last_used": asyncio.get_event_loop().time(),
}
@ -634,9 +650,9 @@ class MCPSessionManager:
self._background_tasks.add(task)
task.add_done_callback(self._background_tasks.discard)
# Wait for session to be ready
# Wait for session to be ready (use longer timeout for remote connections)
try:
session = await asyncio.wait_for(session_future, timeout=10.0)
session = await asyncio.wait_for(session_future, timeout=30.0)
except asyncio.TimeoutError as timeout_err:
# Clean up the failed task
if not task.done():
@ -652,50 +668,136 @@ class MCPSessionManager:
return session, task
async def _create_sse_session(self, session_id: str, connection_params):
"""Create a new SSE session as a background task to avoid context issues."""
async def _create_streamable_http_session(
self, session_id: str, connection_params, preferred_transport: str | None = None
):
"""Create a new Streamable HTTP session with SSE fallback as a background task to avoid context issues.
Args:
session_id: Unique identifier for this session
connection_params: Connection parameters including URL, headers, timeouts
preferred_transport: If set to "sse", skip Streamable HTTP and go directly to SSE
Returns:
tuple: (session, task, transport_used) where transport_used is "streamable_http" or "sse"
"""
import asyncio
from mcp.client.sse import sse_client
from mcp.client.streamable_http import streamablehttp_client
# Create a future to get the session
session_future: asyncio.Future[ClientSession] = asyncio.Future()
# Track which transport succeeded
used_transport: list[str] = []
async def session_task():
"""Background task that keeps the session alive."""
try:
async with sse_client(
connection_params["url"],
connection_params["headers"],
connection_params["timeout_seconds"],
connection_params["sse_read_timeout_seconds"],
) as (read, write):
session = ClientSession(read, write)
async with session:
await session.initialize()
# Signal that session is ready
session_future.set_result(session)
streamable_error = None
# Keep the session alive until cancelled
import anyio
# Skip Streamable HTTP if we know SSE works for this server
if preferred_transport != "sse":
# Try Streamable HTTP first with a quick timeout
try:
await logger.adebug(f"Attempting Streamable HTTP connection for session {session_id}")
# Use a shorter timeout for the initial connection attempt (2 seconds)
async with streamablehttp_client(
url=connection_params["url"],
headers=connection_params["headers"],
timeout=connection_params["timeout_seconds"],
) as (read, write, _):
session = ClientSession(read, write)
async with session:
# Initialize with a timeout to fail fast
await asyncio.wait_for(session.initialize(), timeout=2.0)
used_transport.append("streamable_http")
await logger.ainfo(f"Session {session_id} connected via Streamable HTTP")
# Signal that session is ready
session_future.set_result(session)
event = anyio.Event()
try:
await event.wait()
except asyncio.CancelledError:
await logger.ainfo(f"Session {session_id} is shutting down")
except Exception as e: # noqa: BLE001
if not session_future.done():
session_future.set_exception(e)
# Keep the session alive until cancelled
import anyio
event = anyio.Event()
try:
await event.wait()
except asyncio.CancelledError:
await logger.ainfo(f"Session {session_id} (Streamable HTTP) is shutting down")
except (asyncio.TimeoutError, Exception) as e: # noqa: BLE001
# If Streamable HTTP fails or times out, try SSE as fallback immediately
streamable_error = e
error_type = "timed out" if isinstance(e, asyncio.TimeoutError) else "failed"
await logger.awarning(
f"Streamable HTTP {error_type} for session {session_id}: {e}. Falling back to SSE..."
)
else:
await logger.adebug(f"Skipping Streamable HTTP for session {session_id}, using cached SSE preference")
# Try SSE if Streamable HTTP failed or if SSE is preferred
if streamable_error is not None or preferred_transport == "sse":
try:
await logger.adebug(f"Attempting SSE connection for session {session_id}")
# Extract SSE read timeout from connection params, default to 30s if not present
sse_read_timeout = connection_params.get("sse_read_timeout_seconds", 30)
async with sse_client(
connection_params["url"],
connection_params["headers"],
connection_params["timeout_seconds"],
sse_read_timeout,
) as (read, write):
session = ClientSession(read, write)
async with session:
await session.initialize()
used_transport.append("sse")
fallback_msg = " (fallback)" if streamable_error else " (preferred)"
await logger.ainfo(f"Session {session_id} connected via SSE{fallback_msg}")
# Signal that session is ready
if not session_future.done():
session_future.set_result(session)
# Keep the session alive until cancelled
import anyio
event = anyio.Event()
try:
await event.wait()
except asyncio.CancelledError:
await logger.ainfo(f"Session {session_id} (SSE) is shutting down")
except Exception as sse_error: # noqa: BLE001
# Both transports failed (or just SSE if it was preferred)
if streamable_error:
await logger.aerror(
f"Both Streamable HTTP and SSE failed for session {session_id}. "
f"Streamable HTTP error: {streamable_error}. SSE error: {sse_error}"
)
if not session_future.done():
session_future.set_exception(
ValueError(
f"Failed to connect via Streamable HTTP ({streamable_error}) or SSE ({sse_error})"
)
)
else:
await logger.aerror(f"SSE connection failed for session {session_id}: {sse_error}")
if not session_future.done():
session_future.set_exception(ValueError(f"Failed to connect via SSE: {sse_error}"))
# Start the background task
task = asyncio.create_task(session_task())
self._background_tasks.add(task)
task.add_done_callback(self._background_tasks.discard)
# Wait for session to be ready
# Wait for session to be ready (use longer timeout for remote connections)
try:
session = await asyncio.wait_for(session_future, timeout=10.0)
session = await asyncio.wait_for(session_future, timeout=30.0)
# Log which transport was used
if used_transport:
transport_used = used_transport[0]
await logger.ainfo(f"Session {session_id} successfully established using {transport_used}")
return session, task, transport_used
# This shouldn't happen, but handle it just in case
msg = f"Session {session_id} established but transport not recorded"
raise ValueError(msg)
except asyncio.TimeoutError as timeout_err:
# Clean up the failed task
if not task.done():
@ -705,12 +807,10 @@ class MCPSessionManager:
with contextlib.suppress(asyncio.CancelledError):
await task
self._background_tasks.discard(task)
msg = f"Timeout waiting for SSE session {session_id} to initialize"
msg = f"Timeout waiting for Streamable HTTP/SSE session {session_id} to initialize"
await logger.aerror(msg)
raise ValueError(msg) from timeout_err
return session, task
async def _cleanup_session_by_id(self, server_key: str, session_id: str):
"""Clean up a specific session by server key and session ID."""
if server_key not in self.sessions_by_server:
@ -1056,7 +1156,7 @@ class MCPStdioClient:
await self.disconnect()
class MCPSseClient:
class MCPStreamableHttpClient:
def __init__(self, component_cache=None):
self.session: ClientSession | None = None
self._connection_params = None
@ -1080,67 +1180,15 @@ class MCPSseClient:
self._component_cache.set("mcp_session_manager", session_manager)
return session_manager
async def validate_url(self, url: str | None, headers: dict[str, str] | None = None) -> tuple[bool, str]:
"""Validate the SSE URL before attempting connection."""
async def validate_url(self, url: str | None) -> tuple[bool, str]:
"""Validate the Streamable HTTP URL before attempting connection."""
try:
parsed = urlparse(url)
if not parsed.scheme or not parsed.netloc:
return False, "Invalid URL format. Must include scheme (http/https) and host."
async with httpx.AsyncClient() as client:
try:
# For SSE endpoints, try a GET request with short timeout
# Many SSE servers don't support HEAD requests and return 404
response = await client.get(
url, timeout=2.0, headers={"Accept": "text/event-stream", **(headers or {})}
)
# For SSE, we expect the server to either:
# 1. Start streaming (200)
# 2. Return 404 if HEAD/GET without proper SSE handshake is not supported
# 3. Return other status codes that we should handle gracefully
# Don't fail on 404 since many SSE endpoints return this for non-SSE requests
if response.status_code == HTTP_NOT_FOUND:
# This is likely an SSE endpoint that doesn't support regular GET
# Let the actual SSE connection attempt handle this
return True, ""
# Fail on client errors except 404, but allow server errors and redirects
if (
HTTP_BAD_REQUEST <= response.status_code < HTTP_INTERNAL_SERVER_ERROR
and response.status_code != HTTP_NOT_FOUND
):
return False, f"Server returned client error status: {response.status_code}"
except httpx.TimeoutException:
# Timeout on a short request might indicate the server is trying to stream
# This is actually expected behavior for SSE endpoints
return True, ""
except httpx.NetworkError:
return False, "Network error. Could not reach the server."
else:
return True, ""
except (httpx.HTTPError, ValueError, OSError) as e:
except (ValueError, OSError) as e:
return False, f"URL validation error: {e!s}"
async def pre_check_redirect(self, url: str | None, headers: dict[str, str] | None = None) -> str | None:
"""Check for redirects and return the final URL."""
if url is None:
return url
try:
async with httpx.AsyncClient(follow_redirects=False) as client:
# Use GET with SSE headers instead of HEAD since many SSE servers don't support HEAD
response = await client.get(
url, timeout=2.0, headers={"Accept": "text/event-stream", **(headers or {})}
)
if response.status_code == httpx.codes.TEMPORARY_REDIRECT:
return response.headers.get("Location", url)
# Don't treat 404 as an error here - let the main connection handle it
except (httpx.RequestError, httpx.HTTPError) as e:
await logger.awarning(f"Error checking redirects: {e}")
return url
return True, ""
async def _connect_to_server(
self,
@ -1149,27 +1197,31 @@ class MCPSseClient:
timeout_seconds: int = 30,
sse_read_timeout_seconds: int = 30,
) -> list[StructuredTool]:
"""Connect to MCP server using SSE transport (SDK style)."""
"""Connect to MCP server using Streamable HTTP transport with SSE fallback (SDK style)."""
# Validate and sanitize headers early
validated_headers = _process_headers(headers)
if url is None:
msg = "URL is required for SSE mode"
raise ValueError(msg)
is_valid, error_msg = await self.validate_url(url, validated_headers)
if not is_valid:
msg = f"Invalid SSE URL ({url}): {error_msg}"
msg = "URL is required for StreamableHTTP or SSE mode"
raise ValueError(msg)
url = await self.pre_check_redirect(url, validated_headers)
# Store connection parameters for later use in run_tool
self._connection_params = {
"url": url,
"headers": validated_headers,
"timeout_seconds": timeout_seconds,
"sse_read_timeout_seconds": sse_read_timeout_seconds,
}
# Only validate URL if we don't have a cached session
# This avoids expensive HTTP validation calls when reusing sessions
if not self._connected or not self._connection_params:
is_valid, error_msg = await self.validate_url(url)
if not is_valid:
msg = f"Invalid Streamable HTTP or SSE URL ({url}): {error_msg}"
raise ValueError(msg)
# Store connection parameters for later use in run_tool
# Include SSE read timeout for fallback
self._connection_params = {
"url": url,
"headers": validated_headers,
"timeout_seconds": timeout_seconds,
"sse_read_timeout_seconds": sse_read_timeout_seconds,
}
elif headers:
self._connection_params["headers"] = validated_headers
# If no session context is set, create a default one
if not self._session_context:
@ -1177,18 +1229,21 @@ class MCPSseClient:
import uuid
param_hash = uuid.uuid4().hex[:8]
self._session_context = f"default_sse_{param_hash}"
self._session_context = f"default_http_{param_hash}"
# Get or create a persistent session
# Get or create a persistent session (will try Streamable HTTP, then SSE fallback)
session = await self._get_or_create_session()
response = await session.list_tools()
self._connected = True
return response.tools
async def connect_to_server(self, url: str, headers: dict[str, str] | None = None) -> list[StructuredTool]:
"""Connect to MCP server using SSE transport (SDK style)."""
async def connect_to_server(
self, url: str, headers: dict[str, str] | None = None, sse_read_timeout_seconds: int = 30
) -> list[StructuredTool]:
"""Connect to MCP server using Streamable HTTP with SSE fallback transport (SDK style)."""
return await asyncio.wait_for(
self._connect_to_server(url, headers), timeout=get_settings_service().settings.mcp_server_timeout
self._connect_to_server(url, headers, sse_read_timeout_seconds=sse_read_timeout_seconds),
timeout=get_settings_service().settings.mcp_server_timeout,
)
def set_session_context(self, context_id: str):
@ -1204,12 +1259,14 @@ class MCPSseClient:
# Use cached session manager to get/create persistent session
session_manager = self._get_session_manager()
# Cache session so we can access server-assigned session_id later for DELETE
self.session = await session_manager.get_session(self._session_context, self._connection_params, "sse")
self.session = await session_manager.get_session(
self._session_context, self._connection_params, "streamable_http"
)
return self.session
async def _terminate_remote_session(self) -> None:
"""Attempt to explicitly terminate the remote MCP session via HTTP DELETE (best-effort)."""
# Only relevant for SSE transport
# Only relevant for Streamable HTTP or SSE transport
if not self._connection_params or "url" not in self._connection_params:
return
@ -1255,7 +1312,7 @@ class MCPSseClient:
import uuid
param_hash = uuid.uuid4().hex[:8]
self._session_context = f"default_sse_{param_hash}"
self._session_context = f"default_http_{param_hash}"
max_retries = 2
last_error_type = None
@ -1326,7 +1383,7 @@ class MCPSseClient:
await logger.aerror(msg)
# Clean up failed session from cache
if self._session_context and self._component_cache:
cache_key = f"mcp_session_sse_{self._session_context}"
cache_key = f"mcp_session_http_{self._session_context}"
self._component_cache.delete(cache_key)
self._connected = False
raise ValueError(msg) from e
@ -1364,11 +1421,17 @@ class MCPSseClient:
await self.disconnect()
# Backward compatibility: MCPSseClient is now an alias for MCPStreamableHttpClient
# The new client supports both Streamable HTTP and SSE with automatic fallback
MCPSseClient = MCPStreamableHttpClient
async def update_tools(
server_name: str,
server_config: dict,
mcp_stdio_client: MCPStdioClient | None = None,
mcp_sse_client: MCPSseClient | None = None,
mcp_streamable_http_client: MCPStreamableHttpClient | None = None,
mcp_sse_client: MCPStreamableHttpClient | None = None, # Backward compatibility
) -> tuple[str, list[StructuredTool], dict[str, StructuredTool]]:
"""Fetch server config and update available tools."""
if server_config is None:
@ -1377,11 +1440,17 @@ async def update_tools(
return "", [], {}
if mcp_stdio_client is None:
mcp_stdio_client = MCPStdioClient()
if mcp_sse_client is None:
mcp_sse_client = MCPSseClient()
# Backward compatibility: accept mcp_sse_client parameter
if mcp_streamable_http_client is None:
mcp_streamable_http_client = mcp_sse_client if mcp_sse_client is not None else MCPStreamableHttpClient()
# Fetch server config from backend
mode = "Stdio" if "command" in server_config else "SSE" if "url" in server_config else ""
# Determine mode from config, defaulting to Streamable_HTTP if URL present
mode = server_config.get("mode", "")
if not mode:
mode = "Stdio" if "command" in server_config else "Streamable_HTTP" if "url" in server_config else ""
command = server_config.get("command", "")
url = server_config.get("url", "")
tools = []
@ -1394,7 +1463,7 @@ async def update_tools(
raise
# Determine connection type and parameters
client: MCPStdioClient | MCPSseClient | None = None
client: MCPStdioClient | MCPStreamableHttpClient | None = None
if mode == "Stdio":
# Stdio connection
args = server_config.get("args", [])
@ -1402,10 +1471,10 @@ async def update_tools(
full_command = " ".join([command, *args])
tools = await mcp_stdio_client.connect_to_server(full_command, env)
client = mcp_stdio_client
elif mode == "SSE":
# SSE connection
tools = await mcp_sse_client.connect_to_server(url, headers=headers)
client = mcp_sse_client
elif mode in ["Streamable_HTTP", "SSE"]:
# Streamable HTTP connection with SSE fallback
tools = await mcp_streamable_http_client.connect_to_server(url, headers=headers)
client = mcp_streamable_http_client
else:
logger.error(f"Invalid MCP server mode for '{server_name}': {mode}")
return "", [], {}

View File

@ -229,7 +229,7 @@ class LCModelComponent(Component):
system_message_added = True
runnable = prompt | runnable
else:
messages.append(input_value.to_lc_message())
messages.append(input_value.to_lc_message(self.name))
else:
messages.append(HumanMessage(content=input_value))

View File

@ -7,7 +7,12 @@ from typing import Any
from langchain_core.tools import StructuredTool # noqa: TC002
from lfx.base.agents.utils import maybe_unflatten_dict, safe_cache_get, safe_cache_set
from lfx.base.mcp.util import MCPSseClient, MCPStdioClient, create_input_schema_from_json_schema, update_tools
from lfx.base.mcp.util import (
MCPStdioClient,
MCPStreamableHttpClient,
create_input_schema_from_json_schema,
update_tools,
)
from lfx.custom.custom_component.component_with_cache import ComponentWithCache
from lfx.inputs.inputs import InputTypes # noqa: TC001
from lfx.io import BoolInput, DropdownInput, McpInput, MessageTextInput, Output
@ -32,7 +37,9 @@ class MCPToolsComponent(ComponentWithCache):
# Initialize clients with access to the component cache
self.stdio_client: MCPStdioClient = MCPStdioClient(component_cache=self._shared_component_cache)
self.sse_client: MCPSseClient = MCPSseClient(component_cache=self._shared_component_cache)
self.streamable_http_client: MCPStreamableHttpClient = MCPStreamableHttpClient(
component_cache=self._shared_component_cache
)
def _ensure_cache_structure(self):
"""Ensure the cache has the required structure."""
@ -207,7 +214,7 @@ class MCPToolsComponent(ComponentWithCache):
server_name=server_name,
server_config=server_config,
mcp_stdio_client=self.stdio_client,
mcp_sse_client=self.sse_client,
mcp_streamable_http_client=self.streamable_http_client,
)
self.tool_names = [tool.name for tool in tool_list if hasattr(tool, "name")]
@ -496,7 +503,7 @@ class MCPToolsComponent(ComponentWithCache):
session_context = self._get_session_context()
if session_context:
self.stdio_client.set_session_context(session_context)
self.sse_client.set_session_context(session_context)
self.streamable_http_client.set_session_context(session_context)
exec_tool = self._tool_cache[self.tool]
tool_args = self.get_inputs_for_all_tools(self.tools)[self.tool]

View File

@ -1548,7 +1548,16 @@ class Component(CustomComponent):
if not message.sender_name:
message.sender_name = MESSAGE_SENDER_NAME_AI
async def send_message(self, message: Message, id_: str | None = None):
async def send_message(self, message: Message, id_: str | None = None, *, skip_db_update: bool = False):
"""Send a message with optional database update control.
Args:
message: The message to send
id_: Optional message ID
skip_db_update: If True, only update in-memory and send event, skip DB write.
Useful during streaming to avoid excessive DB round-trips.
Note: This assumes the message already exists in the database with message.id set.
"""
if self._should_skip_message(message):
return message
@ -1558,26 +1567,37 @@ class Component(CustomComponent):
# Ensure required fields for message storage are set
self._ensure_message_required_fields(message)
stored_message = await self._store_message(message)
# If skip_db_update is True and message already has an ID, skip the DB write
# This path is used during agent streaming to avoid excessive DB round-trips
if skip_db_update and message.id:
# Create a fresh Message instance for consistency with normal flow
stored_message = await Message.create(**message.model_dump())
self._stored_message_id = stored_message.id
# Still send the event to update the client in real-time
# Note: If this fails, we don't need DB cleanup since we didn't write to DB
await self._send_message_event(stored_message, id_=id_)
else:
# Normal flow: store/update in database
stored_message = await self._store_message(message)
self._stored_message_id = stored_message.id
try:
complete_message = ""
if (
self._should_stream_message(stored_message, message)
and message is not None
and isinstance(message.text, AsyncIterator | Iterator)
):
complete_message = await self._stream_message(message.text, stored_message)
stored_message.text = complete_message
stored_message = await self._update_stored_message(stored_message)
else:
# Only send message event for non-streaming messages
await self._send_message_event(stored_message, id_=id_)
except Exception:
# remove the message from the database
await delete_message(stored_message.id)
raise
self._stored_message_id = stored_message.id
try:
complete_message = ""
if (
self._should_stream_message(stored_message, message)
and message is not None
and isinstance(message.text, AsyncIterator | Iterator)
):
complete_message = await self._stream_message(message.text, stored_message)
stored_message.text = complete_message
stored_message = await self._update_stored_message(stored_message)
else:
# Only send message event for non-streaming messages
await self._send_message_event(stored_message, id_=id_)
except Exception:
# remove the message from the database
await delete_message(stored_message.id)
raise
self.status = stored_message
return stored_message

View File

@ -34,6 +34,7 @@ class SendMessageFunctionType(Protocol):
id_: str | None = None,
*,
allow_markdown: bool = True,
skip_db_update: bool = False,
) -> Message: ...

View File

@ -121,9 +121,13 @@ class Message(Data):
def to_lc_message(
self,
model_name: str | None = None,
) -> BaseMessage:
"""Converts the Data to a BaseMessage.
Args:
model_name: The model name to use for conversion. Optional.
Returns:
BaseMessage: The converted BaseMessage.
"""
@ -139,7 +143,7 @@ class Message(Data):
if self.sender == MESSAGE_SENDER_USER or not self.sender:
if self.files:
contents = [{"type": "text", "text": text}]
file_contents = self.get_file_content_dicts()
file_contents = self.get_file_content_dicts(model_name)
contents.extend(file_contents)
human_message = HumanMessage(content=contents)
else:
@ -197,7 +201,7 @@ class Message(Data):
return value
# Keep this async method for backwards compatibility
def get_file_content_dicts(self):
def get_file_content_dicts(self, model_name: str | None = None):
content_dicts = []
try:
files = get_file_paths(self.files)
@ -209,7 +213,7 @@ class Message(Data):
if isinstance(file, Image):
content_dicts.append(file.to_content_dict())
else:
content_dicts.append(create_image_content_dict(file))
content_dicts.append(create_image_content_dict(file, None, model_name))
return content_dicts
def load_lc_prompt(self):

View File

@ -56,12 +56,15 @@ def create_data_url(image_path: str | Path, mime_type: str | None = None) -> str
@lru_cache(maxsize=50)
def create_image_content_dict(image_path: str | Path, mime_type: str | None = None) -> dict:
def create_image_content_dict(
image_path: str | Path, mime_type: str | None = None, model_name: str | None = None
) -> dict:
"""Create a content dictionary for multimodal inputs from an image file.
Args:
image_path: Path to the image file
mime_type: MIME type of the image. If None, will be auto-detected
model_name: Optional model parameter to determine content dict structure
Returns:
Content dictionary with type and image_url fields
@ -70,4 +73,7 @@ def create_image_content_dict(image_path: str | Path, mime_type: str | None = No
FileNotFoundError: If the image file doesn't exist
"""
data_url = create_data_url(image_path, mime_type)
if model_name == "OllamaModel":
return {"type": "image_url", "source_type": "url", "image_url": data_url}
return {"type": "image", "source_type": "url", "url": data_url}