Merge remote-tracking branch 'origin/master' into feat/workspace-operation-traceability

This commit is contained in:
TyperBody
2026-09-28 01:56:03 +08:00
25 changed files with 813 additions and 92 deletions
+2
View File
@@ -46,6 +46,7 @@ Run the narrowest useful test first, then broader checks when confidence is need
- Dev environment guide: https://langbot.app/docs/zh/develop/dev-config.
- Plugin runtime / CLI / SDK debugging: https://langbot.app/docs/zh/develop/plugin-runtime.
- API-key auth: `docs/API_KEY_AUTH.md`.
- Service API CLI: [`langbot-cli`](https://github.com/langbot-app/langbot-cli) (`lbctl`) manages running Workspaces; it is separate from the SDK's `lbp` CLI.
- Box deep-dive notes: `docs/review/box-architecture.md` and related files.
- In-repo skills: `skills/` is the single source of truth for LangBot agent skills.
- SDK repo: `../langbot-plugin-sdk/` when changing shared entities, plugin APIs, action protocol, `lbp rt`, or `lbp box`.
@@ -83,6 +84,7 @@ Config keys to verify in `data/config.yaml` / `src/langbot/templates/config.yaml
## Change Rules
- HTTP API changes that should be agent-accessible must update the matching MCP tool in `src/langbot/pkg/api/mcp/server.py` and the relevant skill under `skills/` in the same pass.
- When changing Service API routes or capabilities used by `lbctl`, check compatibility with the separate `langbot-cli` repository.
- New schema changes use Alembic under `src/langbot/pkg/persistence/alembic/versions/`. LangBot 4.x does not support upgrading 3.x databases.
- New platform behavior belongs in platform adapters only for platform translation; pipeline/business logic belongs in `pkg/pipeline/` or services.
- User-facing strings must support i18n (`en_US`, `zh_Hans`; include `ja_JP` where the repo already does).
+2
View File
@@ -23,6 +23,7 @@ LangBot is not a single-repo system.
- `LangBot/` is the main product: backend, web UI, platform adapters, pipeline engine, HTTP API, MCP server, RAG, persistence, skills integration, and the bridge code that talks to runtimes.
- `langbot-plugin-sdk/` is published as `langbot-plugin` and pinned in `LangBot/pyproject.toml`. It contains plugin developer APIs, shared entities, `lbp`, the Plugin Runtime (`lbp rt`), and the Box Runtime (`lbp box`).
- [`langbot-cli`](https://github.com/langbot-app/langbot-cli) is a separately released `lbctl` client for managing a running Workspace through the HTTP Service API; it is not part of the LangBot server or the SDK's `lbp` CLI.
- Plugins import SDK APIs from `langbot_plugin.*`; the LangBot main process imports the same package for shared entities and runtime protocols.
This split matters. If a change modifies SDK entities, component APIs, action protocols, `lbp rt`, or `lbp box`, verify the sibling SDK repo and install the local SDK into LangBot's virtualenv when testing cross-repo behavior.
@@ -252,6 +253,7 @@ LangBot is deliberately agent-friendly. The agent-facing surfaces are part of th
- `skills/` is the single source of truth for in-repo skills.
- `pkg/api/mcp/server.py` exposes the LangBot MCP server at `/mcp`.
- `lbctl` calls the authenticated HTTP Service API from a terminal; its source and releases live in `langbot-cli`.
- `api.global_api_key` authenticates API/MCP access without a browser login.
- `AGENTS.md` and `ARCHITECTURE.md` tell coding agents how the repo works.
+1
View File
@@ -164,6 +164,7 @@ docker compose --profile all up -d
LangBot is **agent-friendly by design** — your coding agents (Claude Code, Codex, Copilot, Cursor, …) can operate, extend, and deploy LangBot with first-class support:
- **MCP Server** — LangBot exposes a built-in [Model Context Protocol](https://modelcontextprotocol.io/) endpoint at `/mcp`, mirroring the HTTP API so an agent can manage bots, pipelines, plugins, and models programmatically. Authenticate with the same API key (set a global key in `config.yaml` or use a per-user key) — no login flow required. Configure it in the Web panel's **API & MCP** tab.
- **CLI (`lbctl`)** — The standalone [LangBot CLI](https://github.com/langbot-app/langbot-cli) lets agents manage a running LangBot Workspace from a terminal through the Service API. It uses an API key and supports bots, pipelines, knowledge bases, models, plugins, skills, and MCP servers. See the CLI repository for installation and commands.
- **In-repo Skills** — The [`skills/`](skills/) directory is the **single source of truth** for working with LangBot: plugin development, core development, end-to-end testing, deployment, and operating the LangBot / LangBot Space MCP servers. Point your agent at this directory and it knows how to build.
- **AGENTS.md** — Every repo ships an [`AGENTS.md`](AGENTS.md) (symlinked to `CLAUDE.md`) describing architecture, conventions, and the rule that API changes must keep the MCP server and skills in sync.
- **`llms.txt`** — Machine-readable project context for LLMs is published on the website.
+1
View File
@@ -180,6 +180,7 @@ docker compose --profile all up -d
LangBot **从设计上就对 Agent 友好** —— 你的编码 Agent(Claude Code、Codex、Copilot、Cursor 等)可以一等公民般地操作、扩展和部署 LangBot:
- **MCP Server** —— LangBot 内置 [Model Context Protocol](https://modelcontextprotocol.io/) 端点 `/mcp`,与 HTTP API 对齐,Agent 可编程式管理机器人、流水线、插件和模型。使用同一套 API Key 鉴权(可在 `config.yaml` 配置全局 Key,或使用用户 Key),无需登录流程。在 Web 面板的 **API 与 MCP** 标签页中配置。
- **CLI(`lbctl`)** —— 独立的 [LangBot CLI](https://github.com/langbot-app/langbot-cli) 让 Agent 在终端通过 Service API 管理已运行实例中的 Workspace。它使用 API Key,支持管理机器人、流水线、知识库、模型、插件、Skill 和 MCP Server。安装方法和命令以 CLI 仓库为准。
- **仓库内 Skills** —— [`skills/`](skills/) 目录是使用 LangBot 的**唯一事实来源**:插件开发、核心开发、端到端测试、部署,以及操作 LangBot / LangBot Space MCP Server。把 Agent 指向这个目录,它就知道如何动手。
- **AGENTS.md** —— 每个仓库都提供 [`AGENTS.md`](AGENTS.md)(软链到 `CLAUDE.md`),描述架构、规范,以及「API 变更必须同步更新 MCP Server 和 skills」的约定。
- **`llms.txt`** —— 面向 LLM 的机器可读项目上下文已发布在官网。
+5
View File
@@ -55,6 +55,11 @@ Behavior:
## Using API Keys
The standalone [`lbctl` CLI](https://github.com/langbot-app/langbot-cli) uses
these API keys to manage a running LangBot Workspace through the Service API.
It can check the key's Workspace identity and server capabilities before
managing resources. See the CLI repository for installation and commands.
### Authentication Headers
Include your API key in the request header using one of these methods:
+1 -1
View File
@@ -71,7 +71,7 @@ dependencies = [
"langchain-text-splitters>=1.1.2",
"chromadb>=1.0.0,<2.0.0",
"qdrant-client (>=1.15.1,<2.0.0)",
"langbot-plugin==0.6.11",
"langbot-plugin==0.6.20",
"asyncpg>=0.30.0",
"line-bot-sdk>=3.19.0",
"matrix-nio>=0.25.2",
+4 -1
View File
@@ -94,7 +94,9 @@ discovered by `importutil.import_modules_in_pkg`.
3. **If the endpoint should be agent-accessible, add/adjust the matching MCP tool
in `pkg/api/mcp/server.py` and update the `langbot-mcp-ops` skill.** API and
MCP surface must stay aligned (see `AGENTS.md`).
4. Update `docs/service-api-openapi.json` if you maintain the OpenAPI overview.
4. If `lbctl` uses the route or capability, check the separate
[`langbot-cli`](https://github.com/langbot-app/langbot-cli) client for compatibility.
5. Update `docs/service-api-openapi.json` if you maintain the OpenAPI overview.
## Database migrations (Alembic)
@@ -127,3 +129,4 @@ uv run python tests/manual/mcp_smoke.py # MCP server e2e smoke
- `langbot-testing` — WebUI/e2e QA harness (`bin/lbs`).
- `langbot-deploy` — Docker/compose deployment + config.
- `langbot-mcp-ops` — operating the LangBot MCP server.
- [`langbot-cli`](https://github.com/langbot-app/langbot-cli) — `lbctl`, a standalone Service API client for managing running Workspaces.
+5
View File
@@ -76,6 +76,11 @@ The tools wrap the LangBot service layer. Current tools (v1):
| `list_knowledge_bases` / `get_knowledge_base` / `retrieve_knowledge_base` | RAG knowledge bases (incl. semantic search) |
| `list_mcp_servers` | External MCP servers LangBot connects to (as a client) |
| `list_skills` / `get_skill` | Installed skills |
| `list_knowledge_engines` / `get_knowledge_engine_schema` / `list_knowledge_parsers` | Discover RAG configuration |
| `get_pipeline_extensions` / `update_pipeline_extensions` | Read or completely replace extension bindings; all lists and switches required |
| `run_pipeline` | One fresh-session turn; requires `runtime.operate`, executes configured models/tools, never auto-retry an unknown outcome |
| `get_monitoring_records` / `get_monitoring_details` | Bounded Workspace records and existing message/session details |
| `get_sandbox_diagnostics` | Read status (`resource.view`), sessions/errors (`audit.view`); managed sandbox admission still applies |
Mutating tools (`create_*`, `update_*`) take a JSON object matching the same
shape as the corresponding HTTP API request body. Discover resources with the
@@ -16,7 +16,7 @@ class BoxRouterGroup(group.RouterGroup):
@self.route(
'/status',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN,
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def _(request_context: RequestContext) -> str:
@@ -42,7 +42,7 @@ class BoxRouterGroup(group.RouterGroup):
@self.route(
'/sessions',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN,
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.AUDIT_VIEW,
)
async def _(request_context: RequestContext) -> str:
@@ -55,7 +55,7 @@ class BoxRouterGroup(group.RouterGroup):
@self.route(
'/errors',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN,
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.AUDIT_VIEW,
)
async def _(request_context: RequestContext) -> str:
@@ -75,7 +75,12 @@ class MonitoringRouterGroup(group.RouterGroup):
return self.success(data=stats)
@self.route('/messages', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/messages',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_messages(request_context: RequestContext) -> str:
"""Get message logs"""
# Parse query parameters
@@ -84,8 +89,9 @@ class MonitoringRouterGroup(group.RouterGroup):
session_ids = quart.request.args.getlist('sessionId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
# Parse datetime
start_time = parse_iso_datetime(start_time_str)
@@ -111,7 +117,12 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/llm-calls', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/llm-calls',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_llm_calls(request_context: RequestContext) -> str:
"""Get LLM call records"""
# Parse query parameters
@@ -119,8 +130,9 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids = quart.request.args.getlist('pipelineId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
# Parse datetime
start_time = parse_iso_datetime(start_time_str)
@@ -145,7 +157,12 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/tool-calls', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/tool-calls',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_tool_calls(request_context: RequestContext) -> str:
"""Get tool call records"""
bot_ids = quart.request.args.getlist('botId')
@@ -153,8 +170,9 @@ class MonitoringRouterGroup(group.RouterGroup):
session_ids = quart.request.args.getlist('sessionId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
start_time = parse_iso_datetime(start_time_str)
end_time = parse_iso_datetime(end_time_str)
@@ -179,15 +197,21 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/embedding-calls', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/embedding-calls',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_embedding_calls(request_context: RequestContext) -> str:
"""Get embedding call records"""
# Parse query parameters
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
knowledge_base_id = quart.request.args.get('knowledgeBaseId')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
# Parse datetime
start_time = parse_iso_datetime(start_time_str)
@@ -211,7 +235,12 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/sessions', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/sessions',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_sessions(request_context: RequestContext) -> str:
"""Get session information"""
# Parse query parameters
@@ -221,8 +250,9 @@ class MonitoringRouterGroup(group.RouterGroup):
end_time_str = quart.request.args.get('endTime')
user_query = quart.request.args.get('userQuery')
is_active_str = quart.request.args.get('isActive')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
# Parse datetime
start_time = parse_iso_datetime(start_time_str)
@@ -254,7 +284,12 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/errors', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/errors',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_errors(request_context: RequestContext) -> str:
"""Get error logs"""
# Parse query parameters
@@ -262,8 +297,9 @@ class MonitoringRouterGroup(group.RouterGroup):
pipeline_ids = quart.request.args.getlist('pipelineId')
start_time_str = quart.request.args.get('startTime')
end_time_str = quart.request.args.get('endTime')
limit = int(quart.request.args.get('limit', 100))
offset = int(quart.request.args.get('offset', 0))
limit, offset = self.ap.monitoring_service.normalize_page_window(
quart.request.args.get('limit', 100), quart.request.args.get('offset', 0)
)
# Parse datetime
start_time = parse_iso_datetime(start_time_str)
@@ -404,7 +440,12 @@ class MonitoringRouterGroup(group.RouterGroup):
}
)
@self.route('/sessions/<session_id>/analysis', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/sessions/<session_id>/analysis',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_session_analysis(session_id: str, request_context: RequestContext) -> str:
"""Get detailed analysis for a specific session"""
start_time = parse_iso_datetime(quart.request.args.get('startTime'))
@@ -421,13 +462,18 @@ class MonitoringRouterGroup(group.RouterGroup):
# The frontend will handle the 'found: false' case
return self.success(data=analysis)
@self.route('/messages/<message_id>/details', methods=['GET'], permission=Permission.RESOURCE_VIEW)
@self.route(
'/messages/<message_id>/details',
methods=['GET'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RESOURCE_VIEW,
)
async def get_message_details(message_id: str, request_context: RequestContext) -> str:
"""Get detailed information for a specific message"""
details = await self.ap.monitoring_service.get_message_details(request_context, message_id)
if not details.get('found'):
return self.error(message=f'Message {message_id} not found', code=404)
return self.http_status(404, 'resource_not_found', 'Message not found')
return self.success(data=details)
@@ -6,6 +6,7 @@ from ....authz import Permission, has_permission
from ....context import RequestContext
from ......operation_trace import service as settings_service
from ....service.secrets import redact_secrets
from ....service.pipeline_run import run_pipeline
from ... import group
from ......pipeline.extension_preferences import (
normalize_extension_preferences,
@@ -248,3 +249,18 @@ class PipelinesRouterGroup(group.RouterGroup):
except ValueError as exc:
return self.http_status(400, -1, str(exc))
return self.success()
@self.route(
'/<pipeline_uuid>/run',
methods=['POST'],
auth_type=group.AuthType.USER_TOKEN_OR_API_KEY,
permission=Permission.RUNTIME_OPERATE,
)
async def run(pipeline_uuid: str, request_context: RequestContext):
body = await quart.request.get_json()
if not isinstance(body, dict) or set(body) != {'message'}:
return self.http_status(400, -1, 'Expected an object containing message')
text = body['message']
if not isinstance(text, str) or not text.strip() or len(text) > 100_000:
return self.http_status(400, -1, 'message must contain 1..100000 characters')
return self.success(data=await run_pipeline(self.ap, request_context, pipeline_uuid, text))
@@ -13,6 +13,24 @@ from .....workspace.invitation_delivery import InvitationDeliveryService
SYSTEM_CAPABILITY_OPERATIONS = (
'knowledge_engine.list',
'knowledge_engine.creation_schema',
'knowledge_engine.retrieval_schema',
'knowledge_parser.list',
'pipeline.extensions.get',
'pipeline.extensions.update',
'pipeline.run',
'monitoring.messages',
'monitoring.llm_calls',
'monitoring.tool_calls',
'monitoring.embedding_calls',
'monitoring.sessions',
'monitoring.errors',
'monitoring.message_details',
'monitoring.session_analysis',
'sandbox.status',
'sandbox.sessions',
'sandbox.errors',
'bot.list',
'bot.get',
'bot.create',
+18 -42
View File
@@ -434,58 +434,34 @@ class KnowledgeService:
# ================= Knowledge Engine Discovery =================
async def list_knowledge_engines(self, context: TenantContext) -> list[dict]:
"""List all available Knowledge Engines from plugins."""
"""List engines; unavailable runtimes must not look like empty catalogs."""
require_workspace_uuid(context)
engines = []
if not self.ap.plugin_connector.is_enable_plugin:
return engines
return []
await self.ap.plugin_connector.require_workspace_context(context)
# Get KnowledgeEngine plugins
try:
knowledge_engines = await self.ap.plugin_connector.list_knowledge_engines()
engines.extend(knowledge_engines)
except Exception as e:
self.ap.logger.warning(f'Failed to list Knowledge Engines from plugins: {e}')
return engines
return await self.ap.plugin_connector.list_knowledge_engines()
async def list_parsers(self, context: TenantContext, mime_type: str | None = None) -> list[dict]:
"""List available parsers, optionally filtered by MIME type."""
"""List parsers, optionally filtered by MIME type."""
require_workspace_uuid(context)
if not self.ap.plugin_connector.is_enable_plugin:
return []
await self.ap.plugin_connector.require_workspace_context(context)
try:
parsers = await self.ap.plugin_connector.list_parsers()
if mime_type:
parsers = [p for p in parsers if mime_type in p.get('supported_mime_types', [])]
return parsers
except Exception as e:
self.ap.logger.warning(f'Failed to list parsers: {e}')
return []
parsers = await self.ap.plugin_connector.list_parsers()
return [p for p in parsers if mime_type in p.get('supported_mime_types', [])] if mime_type else parsers
async def _require_knowledge_engine(self, context: TenantContext, plugin_id: str) -> None:
engines = await self.list_knowledge_engines(context)
if not any(engine.get('plugin_id') == plugin_id for engine in engines):
raise WorkspaceNotFoundError('Knowledge engine not found')
await self.ap.plugin_connector.require_workspace_context(context)
async def get_engine_creation_schema(self, context: TenantContext, plugin_id: str) -> dict:
"""Get creation settings schema for a specific Knowledge Engine."""
require_workspace_uuid(context)
if not self.ap.plugin_connector.is_enable_plugin:
return {}
await self.ap.plugin_connector.require_workspace_context(context)
try:
return await self.ap.plugin_connector.get_rag_creation_schema(plugin_id)
except Exception as e:
self.ap.logger.warning(f'Failed to get creation schema for {plugin_id}: {e}')
return {}
"""Get an existing engine's creation schema in this Workspace."""
await self._require_knowledge_engine(context, plugin_id)
return await self.ap.plugin_connector.get_rag_creation_schema(plugin_id)
async def get_engine_retrieval_schema(self, context: TenantContext, plugin_id: str) -> dict:
"""Get retrieval settings schema for a specific Knowledge Engine."""
require_workspace_uuid(context)
if not self.ap.plugin_connector.is_enable_plugin:
return {}
await self.ap.plugin_connector.require_workspace_context(context)
try:
return await self.ap.plugin_connector.get_rag_retrieval_schema(plugin_id)
except Exception as e:
self.ap.logger.warning(f'Failed to get retrieval schema for {plugin_id}: {e}')
return {}
"""Get an existing engine's retrieval schema in this Workspace."""
await self._require_knowledge_engine(context, plugin_id)
return await self.ap.plugin_connector.get_rag_retrieval_schema(plugin_id)
@@ -332,6 +332,28 @@ class PipelineService:
return new_uuid
async def get_pipeline_extensions(self, context: TenantContext, pipeline_uuid: str) -> dict:
pipeline = await self.get_pipeline(context, pipeline_uuid)
if pipeline is None:
raise WorkspaceNotFoundError('Pipeline not found')
if self.ap.plugin_connector.is_enable_plugin:
await self.ap.plugin_connector.require_workspace_context(context)
plugins = await self.ap.plugin_connector.list_plugins(component_kinds=['Command', 'EventListener', 'Tool'])
prefs = pipeline.get('extensions_preferences', {})
return {
'enable_all_plugins': prefs.get('enable_all_plugins', True),
'enable_all_mcp_servers': prefs.get('enable_all_mcp_servers', True),
'enable_all_skills': prefs.get('enable_all_skills', True),
'bound_plugins': prefs.get('plugins', []),
'available_plugins': redact_secrets(plugins),
'bound_mcp_servers': prefs.get('mcp_servers', []),
'available_mcp_servers': await self.ap.mcp_service.get_mcp_servers(context, contain_runtime_info=True),
'bound_mcp_resources': prefs.get('mcp_resources', []),
'mcp_resource_agent_read_enabled': prefs.get('mcp_resource_agent_read_enabled', True),
'bound_skills': prefs.get('skills', []),
'available_skills': await self.ap.skill_service.list_skills(context),
}
async def update_pipeline_extensions(
self,
context: TenantContext,
@@ -0,0 +1,68 @@
"""Single-turn diagnostic execution through the normal Pipeline scheduler."""
from __future__ import annotations
import asyncio
import time
import uuid
from langbot_plugin.api.entities.builtin.platform import entities, events, message
from langbot_plugin.api.entities.builtin.provider.session import LauncherTypes
from ..authz import Permission, require_permission
from ..context import ExecutionContext, RequestContext
from ....workspace.errors import WorkspaceNotFoundError
async def run_pipeline(ap, context: RequestContext, pipeline_uuid: str, text: str, *, timeout: float = 60) -> dict:
"""Run in a fresh session; a timed-out request is never submitted again."""
require_permission(context, Permission.RUNTIME_OPERATE)
if not isinstance(text, str) or not text.strip() or len(text) > 100_000:
raise ValueError('message must contain 1..100000 characters')
if await ap.pipeline_service.get_pipeline(context, pipeline_uuid) is None:
raise WorkspaceNotFoundError('Pipeline not found')
bot = await ap.platform_mgr.get_websocket_proxy_bot(context)
adapter = bot.adapter
session_id = str(uuid.uuid4())
launcher_id = f'websocket_{pipeline_uuid}:{session_id}'
chain = message.MessageChain([message.Plain(text=text)])
event = events.FriendMessage(
sender=entities.Friend(id=launcher_id, nickname='CLI', remark='CLI'),
message_chain=chain,
time=time.time(),
)
execution = ExecutionContext.from_request(context, bot_uuid='websocket-proxy-bot', pipeline_uuid=pipeline_uuid)
query = await ap.query_pool.add_query(
bot_uuid='websocket-proxy-bot',
launcher_type=LauncherTypes.PERSON,
launcher_id=launcher_id,
sender_id=launcher_id,
message_event=event,
message_chain=chain,
adapter=adapter,
pipeline_uuid=pipeline_uuid,
variables={'_cli_run_status': 'failed'},
execution_context=execution,
)
completed = asyncio.Event()
# add_query returns without yielding after enqueueing, so the scheduler
# cannot remove this query before its completion notification is attached.
object.__setattr__(query, '_completion_event', completed)
data = {'pipeline_uuid': pipeline_uuid, 'session_id': f'person_{launcher_id}', 'query_id': query.query_uuid}
try:
await asyncio.wait_for(completed.wait(), timeout=timeout)
except asyncio.TimeoutError:
return {**data, 'status': 'unknown', 'error': 'Execution status unknown; do not automatically retry.'}
replies = adapter.get_websocket_messages(pipeline_uuid, 'person', session_id)
replies = [reply for reply in replies if reply.get('role') == 'assistant']
data.update(
status=query.variables['_cli_run_status'],
replies=replies,
reply='\n'.join(reply.get('content', '') for reply in replies),
message_id=query.variables.get('_monitoring_message_id'),
)
if data['status'] != 'completed':
data['error'] = 'Pipeline failed or was dropped; inspect the message and call records.'
return data
+113
View File
@@ -96,6 +96,119 @@ class LangBotMCPServer:
}
return _dump(data)
@mcp.tool(description='List knowledge engines and their capabilities.')
async def list_knowledge_engines() -> str:
return _dump(await ap.knowledge_service.list_knowledge_engines(_authorized(Permission.RESOURCE_VIEW)))
@mcp.tool(description='Get creation or retrieval schema for an author/name knowledge engine.')
async def get_knowledge_engine_schema(plugin_id: str, kind: typing.Literal['creation', 'retrieval']) -> str:
context = _authorized(Permission.RESOURCE_VIEW)
if kind == 'creation':
return _dump(await ap.knowledge_service.get_engine_creation_schema(context, plugin_id))
return _dump(await ap.knowledge_service.get_engine_retrieval_schema(context, plugin_id))
@mcp.tool(description='List document parsers, optionally filtered by MIME type.')
async def list_knowledge_parsers(mime_type: str | None = None) -> str:
return _dump(await ap.knowledge_service.list_parsers(_authorized(Permission.RESOURCE_VIEW), mime_type))
@mcp.tool(description='Read Pipeline extension bindings, enable-all switches and available extensions.')
async def get_pipeline_extensions(pipeline_uuid: str) -> str:
return _dump(
await ap.pipeline_service.get_pipeline_extensions(_authorized(Permission.RESOURCE_VIEW), pipeline_uuid)
)
@mcp.tool(description='Replace all Pipeline extension bindings. All binding lists and switches are required.')
async def update_pipeline_extensions(
pipeline_uuid: str,
bound_plugins: list[dict],
bound_mcp_servers: list[str],
bound_skills: list[str],
bound_mcp_resources: list[dict],
enable_all_plugins: bool,
enable_all_mcp_servers: bool,
enable_all_skills: bool,
mcp_resource_agent_read_enabled: bool,
) -> str:
context = _authorized(Permission.RESOURCE_MANAGE)
require_permission(context, Permission.RESOURCE_VIEW)
await ap.pipeline_service.update_pipeline_extensions(
context,
pipeline_uuid,
bound_plugins,
bound_mcp_servers,
enable_all_plugins,
enable_all_mcp_servers,
bound_skills,
enable_all_skills,
bound_mcp_resources,
mcp_resource_agent_read_enabled,
)
return _dump(await ap.pipeline_service.get_pipeline_extensions(context, pipeline_uuid))
@mcp.tool(
description='Run one Pipeline turn in a fresh session. Calls configured models/tools. Never retry an unknown outcome automatically.'
)
async def run_pipeline(pipeline_uuid: str, message: str) -> str:
from ..http.service.pipeline_run import run_pipeline as execute
return _dump(await execute(ap, _authorized(Permission.RUNTIME_OPERATE), pipeline_uuid, message))
@mcp.tool(
description='Read Workspace sandbox status, sessions or recent errors. Does not create execution sessions.'
)
async def get_sandbox_diagnostics(kind: typing.Literal['status', 'sessions', 'errors'] = 'status') -> str:
context = _authorized(Permission.RESOURCE_VIEW if kind == 'status' else Permission.AUDIT_VIEW)
if kind == 'status':
return _dump(await ap.box_service.get_status(context))
if kind == 'sessions':
return _dump(await ap.box_service.get_sessions(context))
if ap.box_service.managed_admission_required:
await ap.box_service.require_workspace_sandbox(context)
return _dump(ap.box_service.get_recent_errors(context))
@mcp.tool(
description='Read bounded Workspace runtime records. Filters follow the monitoring service: bot_ids, pipeline_ids, session_ids, start_time, end_time, knowledge_base_id, user_query, is_active as supported by the record kind.'
)
async def get_monitoring_records(
kind: typing.Literal['messages', 'llm_calls', 'tool_calls', 'embedding_calls', 'sessions', 'errors'],
limit: int = 100,
offset: int = 0,
filters: dict | None = None,
) -> str:
import datetime
context = _authorized(Permission.RESOURCE_VIEW)
allowed = {'start_time', 'end_time'}
if kind != 'embedding_calls':
allowed |= {'bot_ids', 'pipeline_ids'}
if kind in {'messages', 'tool_calls'}:
allowed.add('session_ids')
if kind == 'embedding_calls':
allowed.add('knowledge_base_id')
if kind == 'sessions':
allowed |= {'user_query', 'is_active'}
kwargs = dict(filters or {})
if set(kwargs) - allowed:
raise ValueError('Unsupported monitoring filter')
for field in ('start_time', 'end_time'):
if kwargs.get(field):
value = datetime.datetime.fromisoformat(kwargs[field].replace('Z', '+00:00'))
if value.tzinfo is not None:
value = value.astimezone(datetime.timezone.utc).replace(tzinfo=None)
kwargs[field] = value
limit, offset = ap.monitoring_service.normalize_page_window(limit, offset)
rows, total = await getattr(ap.monitoring_service, 'get_' + kind)(
context, limit=limit, offset=offset, **kwargs
)
return _dump({kind: rows, 'total': total, 'limit': limit, 'offset': offset})
@mcp.tool(description='Read message details or session analysis within the authenticated Workspace.')
async def get_monitoring_details(kind: typing.Literal['message', 'session'], identifier: str) -> str:
context = _authorized(Permission.RESOURCE_VIEW)
if kind == 'message':
return _dump(await ap.monitoring_service.get_message_details(context, identifier))
return _dump(await ap.monitoring_service.get_session_analysis(context, identifier))
# ----- Bots ---------------------------------------------------- #
@mcp.tool(description='List all messaging-platform bots. Secrets are redacted.')
async def list_bots() -> str:
+6 -5
View File
@@ -1373,13 +1373,14 @@ class BoxService:
self._runtime_connector.dispose()
async def get_sessions(self, context: TenantContext) -> list[dict]:
if not self._available:
if not self._enabled:
if self.managed_admission_required:
await self.require_workspace_sandbox(context)
return []
execution_context = await self.require_workspace_sandbox(context)
try:
return await self.client.get_sessions(action_context=self._action_context(execution_context))
except Exception:
return []
if not self._available:
raise RuntimeError('Sandbox is unavailable; inspect sandbox status')
return await self.client.get_sessions(action_context=self._action_context(execution_context))
async def get_storage_analysis(self, context: TenantContext) -> dict:
"""Return Workspace-scoped storage measured by the Box Runtime."""
+4
View File
@@ -444,6 +444,10 @@ class RuntimePipeline:
self.ap.logger.debug(f'Processing query {query.query_id}')
await self._execute_from_stage(0, query)
if '_cli_run_status' in query.variables:
query.variables['_cli_run_status'] = (
'failed' if query.variables.get('_monitoring_has_error') else 'completed'
)
# Record query success only if no error occurred during processing
if not query.variables.get('_monitoring_has_error', False):
+6
View File
@@ -173,6 +173,9 @@ class QueryPool:
else:
self.active_query_count_by_workspace.pop(query_workspace_uuid, None)
plugin_diagnostics.discard_query_state(query)
completion = getattr(query, '_completion_event', None)
if completion is not None:
completion.set()
counter_key = (
execution_context.instance_uuid,
execution_context.workspace_uuid,
@@ -426,6 +429,9 @@ class QueryPool:
self.queries.pop(index)
break
plugin_diagnostics.discard_query_state(query)
completion = getattr(query, '_completion_event', None)
if completion is not None:
completion.set()
return True
async def __aenter__(self):
+23 -1
View File
@@ -1950,7 +1950,29 @@ class PluginRuntimeConnector(ManagedRuntimeConnector):
if task_context is not None:
task_context.set_current_action('waiting for plugin initialization')
task_context.metadata['progress_percent'] = 84
await self._wait_for_installed_plugin_ready(plugin_author, plugin_name, task_context)
startup_timeout = 30.0
if desired.execution_mode is PluginExecutionMode.SHARED_CERTIFIED:
worker_policy = self.worker_policy or PluginWorkerPolicy(
max_cpus=1.0,
max_memory_mb=512,
max_pids=64,
max_open_files=1024,
max_file_size_mb=128,
)
startup_timeout = (
120.0
* (
math.ceil(worker_policy.max_workers / worker_policy.max_concurrent_restarts)
+ worker_policy.max_installations
)
+ 5.0
)
await self._wait_for_installed_plugin_ready(
plugin_author,
plugin_name,
task_context,
timeout=startup_timeout,
)
if task_context is not None:
task_context.set_current_action('refreshing plugin components')
task_context.metadata['progress_percent'] = 95
+3
View File
@@ -90,6 +90,9 @@ def fake_monitoring_app():
# Monitoring service
app.monitoring_service = Mock()
from langbot.pkg.api.http.service.monitoring import MonitoringService
app.monitoring_service.normalize_page_window = MonitoringService(app).normalize_page_window
app.monitoring_service.get_overview_metrics = AsyncMock(
return_value={
'total_messages': 100,
+18
View File
@@ -576,6 +576,24 @@ async def test_api_key_context_returns_bound_identity_without_workspace_permissi
'mcp_server.resource_read',
'mcp_server.logs',
'mcp_server.test',
'knowledge_engine.list',
'knowledge_engine.creation_schema',
'knowledge_engine.retrieval_schema',
'knowledge_parser.list',
'pipeline.extensions.get',
'pipeline.extensions.update',
'pipeline.run',
'monitoring.messages',
'monitoring.llm_calls',
'monitoring.tool_calls',
'monitoring.embedding_calls',
'monitoring.sessions',
'monitoring.errors',
'monitoring.message_details',
'monitoring.session_analysis',
'sandbox.status',
'sandbox.sessions',
'sandbox.errors',
]
)
assert all(item == {'supported': True} for item in capabilities['operations'].values())
@@ -0,0 +1,384 @@
from __future__ import annotations
import asyncio
from types import SimpleNamespace
from unittest.mock import AsyncMock, Mock
import pytest
import quart
from langbot.pkg.api.http.context import (
ExecutionContext,
PrincipalContext,
PrincipalType,
RequestContext,
WorkspaceContext,
)
from langbot.pkg.api.http.controller.groups.box import BoxRouterGroup
from langbot.pkg.api.http.controller.groups.monitoring import MonitoringRouterGroup
from langbot.pkg.api.http.controller.groups.pipelines.pipelines import PipelinesRouterGroup
from langbot.pkg.api.http.service.pipeline import PipelineService
from langbot.pkg.api.http.service.pipeline_run import run_pipeline
from langbot.pkg.pipeline.pool import QueryPool
from langbot.pkg.pipeline.pipelinemgr import RuntimePipeline, StageInstContainer
from langbot.pkg.pipeline.entities import StageProcessResult, ResultType
from langbot.pkg.platform.sources.websocket_adapter import WebSocketAdapter, WebSocketSession
from langbot.pkg.cloud.entitlements import EntitlementUnavailableError
from .test_monitoring_tenancy import service as service, WORKSPACE_A, WORKSPACE_B, _context, _record_message
pytestmark = pytest.mark.asyncio
def request_context(workspace=WORKSPACE_A, permissions=('resource.view', 'runtime.operate', 'audit.view')):
return RequestContext(
instance_uuid='instance',
placement_generation=3,
request_id='request',
auth_type='api_key',
principal=PrincipalContext(PrincipalType.API_KEY, api_key_uuid=f'key-{workspace}'),
workspace=WorkspaceContext(workspace, None, None, frozenset(permissions)),
)
async def diagnostic_client(ap):
async def authenticate(key):
if key not in {'a', 'b', 'viewer', 'none'}:
raise ValueError('Invalid API key')
ctx = request_context(
WORKSPACE_B if key == 'b' else WORKSPACE_A,
()
if key == 'none'
else ('resource.view',)
if key == 'viewer'
else ('resource.view', 'runtime.operate', 'resource.manage', 'audit.view'),
)
return SimpleNamespace(
instance_uuid=ctx.instance_uuid,
workspace_uuid=ctx.workspace_uuid,
placement_generation=ctx.placement_generation,
api_key_uuid=ctx.principal.api_key_uuid,
permissions=ctx.workspace.permissions,
)
ap.apikey_service = SimpleNamespace(authenticate_api_key=authenticate)
app = quart.Quart(__name__)
for cls in (MonitoringRouterGroup, BoxRouterGroup, PipelinesRouterGroup):
await cls(ap, app).initialize()
return app.test_client()
async def test_http_api_key_diagnostics_are_scoped_bounded_and_permission_checked(service): # noqa: F811
ap = service.ap
ap.logger = Mock()
ap.monitoring_service = service
ap.box_service = SimpleNamespace(
get_status=AsyncMock(return_value={'enabled': False, 'available': False}),
get_sessions=AsyncMock(return_value=[]),
get_recent_errors=Mock(return_value=[]),
managed_admission_required=False,
)
client = await diagnostic_client(ap)
message_a = await _record_message(service, _context(WORKSPACE_A), 'only-a')
await _record_message(service, _context(WORKSPACE_B), 'only-b')
for key, expected, excluded in [('a', 'only-a', 'only-b'), ('b', 'only-b', 'only-a')]:
response = await client.get(
'/api/v1/monitoring/messages?limit=1', headers={'X-API-Key': key, 'X-Workspace-Id': WORKSPACE_B}
)
assert response.status_code == 200
text = await response.get_data(as_text=True)
assert expected in text and excluded not in text
response = await client.get(f'/api/v1/monitoring/messages/{message_a}/details', headers={'X-API-Key': 'b'})
assert response.status_code == 404
assert (await client.get('/api/v1/monitoring/messages', headers={'X-API-Key': 'none'})).status_code == 403
assert (await client.get('/api/v1/monitoring/messages')).status_code == 401
for route in ('llm-calls', 'tool-calls', 'embedding-calls', 'sessions', 'errors'):
assert (await client.get('/api/v1/monitoring/' + route, headers={'X-API-Key': 'a'})).status_code == 200
for route in ('sessions', 'errors'):
assert (await client.get('/api/v1/box/' + route, headers={'X-API-Key': 'viewer'})).status_code == 403
response = await client.get('/api/v1/box/' + route, headers={'X-API-Key': 'a'})
assert response.status_code == 200 and (await response.get_json())['data'] == []
assert (await client.get('/api/v1/box/runtime-status', headers={'X-API-Key': 'a'})).status_code == 401
response = await client.get('/api/v1/box/status', headers={'X-API-Key': 'viewer'})
assert (await response.get_json())['data']['enabled'] is False
ap.box_service.get_status.side_effect = EntitlementUnavailableError('Sandbox unavailable')
response = await client.get('/api/v1/box/status', headers={'X-API-Key': 'a'})
assert response.status_code == 403
assert (await response.get_json())['code'] == 'managed_sandbox_unavailable'
def runtime_app(monitoring_service):
ap = monitoring_service.ap
ap.logger = Mock()
ap.monitoring_service = monitoring_service
ap.query_pool = QueryPool()
ap.workspace_service = SimpleNamespace(
get_execution_binding=AsyncMock(return_value=SimpleNamespace(instance_uuid='instance'))
)
ap.plugin_connector = SimpleNamespace(
emit_event=AsyncMock(
side_effect=lambda event_obj, *_: SimpleNamespace(event=event_obj, is_prevented_default=lambda: False)
)
)
ap.bot_service = SimpleNamespace(get_bot=AsyncMock(return_value={'name': 'CLI'}))
ap.pipeline_service = SimpleNamespace(get_pipeline=AsyncMock(return_value={'uuid': 'p'}))
adapters = {}
async def proxy(context):
if context.workspace_uuid not in adapters:
logger = Mock(execution_context=ExecutionContext.from_request(context))
adapters[context.workspace_uuid] = WebSocketAdapter.model_construct(ap=ap, logger=logger)
adapters[context.workspace_uuid].websocket_person_session = WebSocketSession(id='person')
adapters[context.workspace_uuid].websocket_group_session = WebSocketSession(id='group')
return SimpleNamespace(adapter=adapters[context.workspace_uuid])
ap.platform_mgr = SimpleNamespace(get_websocket_proxy_bot=proxy)
return ap
async def process_next(ap, *, fail=False):
while not ap.query_pool.queries:
await asyncio.sleep(0)
query = ap.query_pool.queries[0]
async with ap.query_pool:
ap.query_pool.mark_query_running_locked(query)
class ReplyStage:
async def process(self, query, _name):
return StageProcessResult(
result_type=ResultType.CONTINUE,
new_query=query,
user_notice='reply:' + str(query.message_chain),
error_notice='model unavailable' if fail else '',
)
pipeline = RuntimePipeline(
ap,
SimpleNamespace(
uuid='p',
workspace_uuid=query.workspace_uuid,
name='Test',
extensions_preferences={},
config={'output': {'misc': {'at-sender': False, 'quote-origin': False}}},
),
[StageInstContainer('reply', ReplyStage())],
ExecutionContext.from_request(request_context(query.workspace_uuid), pipeline_uuid='p'),
)
await pipeline.run(query)
return query
async def test_run_route_executes_pipeline_and_isolates_sessions_and_workspaces(service): # noqa: F811
ap = runtime_app(service)
client = await diagnostic_client(ap)
assert (
await client.post('/api/v1/pipelines/p/run', headers={'X-API-Key': 'viewer'}, json={'message': 'x'})
).status_code == 403
sessions = set()
for key, text in [('a', 'first'), ('a', 'second'), ('b', 'other-workspace')]:
task = asyncio.create_task(process_next(ap))
response = await asyncio.wait_for(
client.post('/api/v1/pipelines/p/run', headers={'X-API-Key': key}, json={'message': text}), 3
)
if response.status_code != 200:
task.cancel()
await asyncio.gather(task, return_exceptions=True)
pytest.fail(await response.get_data(as_text=True))
query = await task
assert response.status_code == 200
data = (await response.get_json())['data']
assert data['status'] == 'completed'
assert data['reply'] == 'reply:' + text
assert data['session_id'] not in sessions
sessions.add(data['session_id'])
assert query.workspace_uuid == (WORKSPACE_B if key == 'b' else WORKSPACE_A)
assert data['message_id']
a_records, _ = await service.get_messages(_context(WORKSPACE_A))
assert 'other-workspace' not in str(a_records)
async def test_run_failure_timeout_and_missing_pipeline(service): # noqa: F811
ap = runtime_app(service)
task = asyncio.create_task(process_next(ap, fail=True))
failed = await run_pipeline(ap, request_context(), 'p', 'fail')
await task
assert failed['status'] == 'failed' and failed['message_id']
errors, total = await service.get_errors(_context(WORKSPACE_A))
assert total == 1 and errors[0]['error_message'] == 'model unavailable'
unknown = await run_pipeline(ap, request_context(), 'p', 'slow', timeout=0.001)
assert unknown['status'] == 'unknown' and len(ap.query_pool.queries) == 1
await process_next(ap)
assert not ap.query_pool.cached_queries
ap.pipeline_service.get_pipeline.return_value = None
from langbot.pkg.workspace.errors import WorkspaceNotFoundError
with pytest.raises(WorkspaceNotFoundError):
await run_pipeline(ap, request_context(), 'missing', 'x')
assert not ap.query_pool.queries
async def test_extension_discovery_preserves_empty_bindings_and_redacts_plugins():
ap = SimpleNamespace(
plugin_connector=SimpleNamespace(
is_enable_plugin=True,
require_workspace_context=AsyncMock(),
list_plugins=AsyncMock(return_value=[{'api_key': 'private'}]),
),
mcp_service=SimpleNamespace(get_mcp_servers=AsyncMock(return_value=[])),
skill_service=SimpleNamespace(list_skills=AsyncMock(return_value=[])),
)
pipeline_service = PipelineService(ap)
pipeline_service.get_pipeline = AsyncMock(
return_value={'extensions_preferences': {'enable_all_plugins': False, 'plugins': []}}
)
data = await pipeline_service.get_pipeline_extensions(request_context(), 'p')
assert data['bound_plugins'] == [] and data['enable_all_plugins'] is False
assert data['available_plugins'] == [{'api_key': '***'}]
ap.plugin_connector.require_workspace_context.assert_awaited_once_with(request_context())
async def test_compiled_cli_against_core_http(service, tmp_path): # noqa: F811
"""Optional cross-repo check: real HTTP/SQLite/Pipeline, synthetic runtime providers."""
import json
import os
import socket
from pathlib import Path
from hypercorn.asyncio import serve
from hypercorn.config import Config
import sqlalchemy
from langbot.pkg.api.http.controller.groups.system import SystemRouterGroup
from langbot.pkg.api.http.controller.groups.knowledge.engines import KnowledgeEnginesRouterGroup
from langbot.pkg.api.http.controller.groups.knowledge.parsers import ParsersRouterGroup
from langbot.pkg.api.http.service.knowledge import KnowledgeService
from langbot.pkg.entity.persistence.pipeline import LegacyPipeline
binary = os.environ.get('LANGBOT_CLI_BIN')
if not binary:
pytest.skip('Set LANGBOT_CLI_BIN to a compiled lbctl for the cross-repo HTTP check')
assert Path(binary).is_file()
ap = runtime_app(service)
ap.pipeline_service = PipelineService(ap)
ap.pipeline_mgr = SimpleNamespace(remove_pipeline=AsyncMock(), load_pipeline=AsyncMock())
await ap.persistence_mgr.execute_async(
sqlalchemy.insert(LegacyPipeline).values(
uuid='p',
workspace_uuid=WORKSPACE_A,
name='CLI test',
description='',
for_version='1',
stages=[],
config={},
)
)
ap.plugin_connector.is_enable_plugin = True
ap.plugin_connector.require_workspace_context = AsyncMock()
ap.plugin_connector.list_plugins = AsyncMock(return_value=[])
ap.plugin_connector.list_knowledge_engines = AsyncMock(
return_value=[{'plugin_id': 'author/engine', 'capabilities': ['doc_ingestion']}]
)
ap.plugin_connector.get_rag_creation_schema = AsyncMock(return_value={'type': 'object'})
ap.plugin_connector.get_rag_retrieval_schema = AsyncMock(return_value={'type': 'object'})
ap.plugin_connector.list_parsers = AsyncMock(return_value=[{'id': 'text', 'supported_mime_types': ['text/plain']}])
ap.knowledge_service = KnowledgeService(ap)
ap.mcp_service = SimpleNamespace(get_mcp_servers=AsyncMock(return_value=[]))
ap.skill_service = SimpleNamespace(list_skills=AsyncMock(return_value=[]))
ap.box_service = SimpleNamespace(
get_status=AsyncMock(return_value={'enabled': False, 'available': False}),
get_sessions=AsyncMock(return_value=[]),
get_recent_errors=Mock(return_value=[]),
managed_admission_required=False,
)
client = await diagnostic_client(ap)
original_authenticate = ap.apikey_service.authenticate_api_key
async def smoke_authenticate(key):
return await original_authenticate('a' if key == 'lbk_synthetic_smoke_key_7f36c0' else 'invalid')
ap.apikey_service.authenticate_api_key = smoke_authenticate
for cls in (SystemRouterGroup, KnowledgeEnginesRouterGroup, ParsersRouterGroup):
await cls(ap, client.app).initialize()
listener = socket.socket()
listener.bind(('127.0.0.1', 0))
port = listener.getsockname()[1]
config = Config()
config.bind = [f'fd://{listener.detach()}']
config.accesslog = None
shutdown = asyncio.Event()
server = asyncio.create_task(serve(client.app, config, shutdown_trigger=shutdown.wait))
cli_config = tmp_path / 'cli.yaml'
env = {k: v for k, v in os.environ.items() if not k.startswith('LANGBOT_')}
env['SMOKE_KEY'] = 'lbk_synthetic_smoke_key_7f36c0'
async def cli(*args, code=0):
process = await asyncio.create_subprocess_exec(
binary,
'--config',
str(cli_config),
'-o',
'json',
*args,
env=env,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.PIPE,
)
out, err = await asyncio.wait_for(process.communicate(), 5)
assert process.returncode == code, (out.decode(), err.decode())
return json.loads(out)
try:
async with asyncio.timeout(3):
while True:
try:
_, writer = await asyncio.open_connection('127.0.0.1', port)
writer.close()
await writer.wait_closed()
break
except OSError:
await asyncio.sleep(0.01)
await cli(
'context',
'add',
'smoke',
'--endpoint',
f'http://127.0.0.1:{port}',
'--api-key-env',
'SMOKE_KEY',
'--expect-workspace',
WORKSPACE_A,
)
await cli('context', 'use', 'smoke')
for args in [
('knowledge-engine', 'list'),
('knowledge-engine', 'creation-schema', 'author/engine'),
('knowledge-engine', 'retrieval-schema', 'author/engine'),
('knowledge-parser', 'list', '--mime-type', 'text/plain'),
('sandbox', 'status'),
('sandbox', 'sessions'),
('sandbox', 'errors'),
]:
assert (await cli(*args))['ok']
await cli('knowledge-engine', 'creation-schema', 'missing/engine', code=5)
body = {
'bound_plugins': [],
'bound_mcp_servers': [],
'bound_skills': [],
'bound_mcp_resources': [],
'enable_all_plugins': False,
'enable_all_mcp_servers': False,
'enable_all_skills': False,
'mcp_resource_agent_read_enabled': False,
}
body_path = tmp_path / 'extensions.json'
body_path.write_text(json.dumps(body))
updated = await cli('pipeline', 'extensions', 'update', 'p', '--file', str(body_path))
assert all(updated['data'][k] == v for k, v in body.items())
task = asyncio.create_task(process_next(ap))
result = await cli('pipeline', 'run', 'p', '--message', 'HTTP smoke')
await task
assert result['data']['status'] == 'completed' and result['data']['reply'] == 'reply:HTTP smoke'
await cli('monitoring', 'message', result['data']['message_id'])
await cli('monitoring', 'session', result['data']['session_id'])
for kind in ('messages', 'llm-calls', 'tool-calls', 'embedding-calls', 'sessions', 'errors'):
assert (await cli('monitoring', kind, '--limit', '1'))['ok']
finally:
shutdown.set()
await server
@@ -53,7 +53,7 @@ def _app():
require_workspace_context=AsyncMock(side_effect=lambda context: context),
get_rag_creation_schema=AsyncMock(return_value={}),
get_rag_retrieval_schema=AsyncMock(return_value={}),
list_knowledge_engines=AsyncMock(return_value=[]),
list_knowledge_engines=AsyncMock(return_value=[{'plugin_id': 'author/engine'}]),
list_parsers=AsyncMock(return_value=[]),
),
)
@@ -338,14 +338,15 @@ async def test_schema_validation_refences_before_second_runtime_call():
@pytest.mark.asyncio
async def test_engine_schemas_are_context_gated_and_fail_soft_on_connector_error():
async def test_engine_schemas_are_context_gated_and_report_connector_error():
app = _app()
app.plugin_connector.get_rag_creation_schema.return_value = {'schema': ['creation']}
app.plugin_connector.get_rag_retrieval_schema.side_effect = RuntimeError('offline')
service = KnowledgeService(app)
assert await service.get_engine_creation_schema(CONTEXT, 'author/engine') == {'schema': ['creation']}
assert await service.get_engine_retrieval_schema(CONTEXT, 'author/engine') == {}
with pytest.raises(RuntimeError, match='offline'):
await service.get_engine_retrieval_schema(CONTEXT, 'author/engine')
with pytest.raises(WorkspaceRequiredError):
await service.get_engine_creation_schema(None, 'author/engine')
@@ -535,12 +536,12 @@ class TestListKnowledgeEngines:
app.plugin_connector.list_knowledge_engines.assert_not_awaited()
@pytest.mark.asyncio
async def test_returns_empty_on_exception(self):
async def test_reports_runtime_exception(self):
app = _app()
app.plugin_connector.list_knowledge_engines.side_effect = RuntimeError('Connection error')
assert await KnowledgeService(app).list_knowledge_engines(CONTEXT) == []
app.logger.warning.assert_called_once()
with pytest.raises(RuntimeError, match='Connection error'):
await KnowledgeService(app).list_knowledge_engines(CONTEXT)
class TestListParsers:
@@ -603,17 +604,12 @@ class TestGetEngineSchemas:
assert 'properties' in result
@pytest.mark.asyncio
async def test_returns_empty_dict_on_exception(self):
async def test_reports_schema_runtime_exception(self):
app = _app()
app.plugin_connector.get_rag_creation_schema.side_effect = RuntimeError('Plugin error')
result = await KnowledgeService(app).get_engine_creation_schema(
CONTEXT,
'author/engine',
)
assert result == {}
app.logger.warning.assert_called_once()
with pytest.raises(RuntimeError, match='Plugin error'):
await KnowledgeService(app).get_engine_creation_schema(CONTEXT, 'author/engine')
class TestKnowledgeBaseSecretViews:
@@ -652,3 +648,12 @@ class TestKnowledgeBaseSecretViews:
)
app.rag_mgr.create_knowledge_base.assert_not_awaited()
@pytest.mark.asyncio
async def test_missing_engine_is_not_an_empty_schema():
app = _app()
app.plugin_connector.list_knowledge_engines.return_value = []
with pytest.raises(WorkspaceNotFoundError, match='Knowledge engine not found'):
await KnowledgeService(app).get_engine_creation_schema(CONTEXT, 'missing/engine')
app.plugin_connector.get_rag_creation_schema.assert_not_awaited()
Generated
+4 -4
View File
@@ -2185,7 +2185,7 @@ requires-dist = [
{ name = "gewechat-client", specifier = ">=0.1.5" },
{ name = "html2text", specifier = ">=2024.2.26" },
{ name = "httpx", extras = ["socks"], specifier = ">=0.28.1" },
{ name = "langbot-plugin", specifier = "==0.6.11" },
{ name = "langbot-plugin", specifier = "==0.6.20" },
{ name = "langchain", specifier = ">=1.3.9" },
{ name = "langchain-core", specifier = ">=1.3.3" },
{ name = "langchain-text-splitters", specifier = ">=1.1.2" },
@@ -2255,7 +2255,7 @@ dev = [
[[package]]
name = "langbot-plugin"
version = "0.6.11"
version = "0.6.20"
source = { registry = "https://pypi.org/simple" }
dependencies = [
{ name = "aiofiles" },
@@ -2276,9 +2276,9 @@ dependencies = [
{ name = "watchdog" },
{ name = "websockets" },
]
sdist = { url = "https://files.pythonhosted.org/packages/41/13/7f913f88494205abb1a6390746e88de4e1148fbf6d0a273bbe2450f5fb0c/langbot_plugin-0.6.11.tar.gz", hash = "sha256:2ecddc2ef3a2eb92925c08cf0f19887e9be573b93e7d5de49494c01c083cfac7", size = 655797, upload-time = "2026-09-26T18:15:19.81Z" }
sdist = { url = "https://files.pythonhosted.org/packages/dd/bb/08bfa37c29a9b6823678e0b837adddacab162b982b36737e0b5e7bea6c40/langbot_plugin-0.6.20.tar.gz", hash = "sha256:86b3803fb83e67379db4a6cf166dec67a8233fbfbc1dbbf744b2b765d0615879", size = 663741, upload-time = "2026-09-27T08:56:36.473Z" }
wheels = [
{ url = "https://files.pythonhosted.org/packages/25/10/00e3d5d985d82aa08dd5f2f01aa2b85e8b46880dcfa2aa9b923f7be285da/langbot_plugin-0.6.11-py3-none-any.whl", hash = "sha256:85458919c943894780db85bbd96c0a61354e231bc396a3f74072a8eddc978b07", size = 427659, upload-time = "2026-09-26T18:15:18.003Z" },
{ url = "https://files.pythonhosted.org/packages/48/28/f537c212f67afd3820bb56920b9b29783b68096fb6b2ca1a388d268dd27c/langbot_plugin-0.6.20-py3-none-any.whl", hash = "sha256:b110878d880b24ac739b26b3eb158067fb4d0be8b55ba3baf97d2092da920f10", size = 429929, upload-time = "2026-09-27T08:56:34.995Z" },
]
[[package]]