mirror of
https://github.com/TencentCloud/Octop.git
synced 2026-10-02 07:34:38 +08:00
feat(connectors): 企查查改为 API Key,并用 internal 模式加载工具
公网 HTTP 无法完成 OAuth 回调;对话若按 gateway 注入会得到空工具表。改为粘贴 Key,harness 走内部 HTTP 聚合五个 MCP。 Co-authored-by: Cursor <cursoragent@cursor.com>
This commit is contained in:
committed by
jubaoliang
co-authored by
Cursor
parent
6f7815bb71
commit
df7f7d187c
+3
-1
@@ -8,10 +8,12 @@
|
||||
|
||||
### 新增
|
||||
|
||||
- 内置企查查 MCP 连接器,一次 OAuth 授权接入企业、风险、知识产权、经营及董监高五类数据,共享刷新与解绑
|
||||
- 内置企查查 MCP 连接器:粘贴一份 API Key 接入企业、风险、知识产权、经营及董监高五类数据;未开通的类别会跳过
|
||||
- 内置 OpenAlex MCP 连接器:OAuth 登录后按用户自己的 API Key 与每日预算检索学术文献、引文与研究实体
|
||||
|
||||
### 修复
|
||||
|
||||
- 企查查在对话中改为按内部 HTTP MCP 加载工具(`mcp_mode=internal`),不再误走进程内 gateway 导致「无法加载 MCP 工具」
|
||||
- 定时任务以专家方式执行时,运行失败(工具调用报错、需要人工介入、没有可见回复)也会把任务提示词和已产出的部分内容投影进会话线程;此前这些线程一个字都没有,点「立即执行」后打开对话只看到空会话(#516)
|
||||
|
||||
## [1.0.2b1] - 2026-09-22
|
||||
|
||||
@@ -18,7 +18,7 @@ export interface ConnectorCatalogEntry {
|
||||
icon: string;
|
||||
color: string;
|
||||
phase: "available" | "coming_soon";
|
||||
mcp_mode: "remote" | "gateway";
|
||||
mcp_mode: "remote" | "gateway" | "internal";
|
||||
category: ConnectorCategory;
|
||||
quick_auth_url?: string | null;
|
||||
login_url?: string | null;
|
||||
|
||||
+13
-36
@@ -1,6 +1,6 @@
|
||||
# 企查查 MCP 连接器
|
||||
|
||||
一张企查查卡片、一次 OAuth 登录,连接以下五个企业数据 Server。查询范围与额度以企查查账户权限为准。
|
||||
一张企查查卡片、一份 API Key,连接以下五个企业数据 Server。查询范围与额度以企查查账户权限为准。
|
||||
|
||||
| 数据类别 | MCP 地址 | 工具名称前缀 |
|
||||
| --- | --- | --- |
|
||||
@@ -13,52 +13,29 @@
|
||||
## 使用
|
||||
|
||||
1. 打开「连接器 → 内置连接器」,选择「企查查」。
|
||||
2. 点击「一键授权」,在企查查页面登录并完成授权,无需手动填写 API Key。
|
||||
3. 返回 Octop,确认五类服务探测成功,保存连接器,并在对话中选择它。
|
||||
2. 登录 [企查查智能体平台](https://agent.qcc.com/),在个人中心复制 API Key,粘贴到 Octop。无需 OAuth 回调。
|
||||
3. 确认已开通的类别探测成功(未开通的类别会跳过),保存连接器,并在对话中选择它。
|
||||
4. 输入企业全称及查询需求,例如「查询思必驰科技股份有限公司的工商变更」。
|
||||
|
||||
授权以 company 为入口,复用动态客户端注册、PKCE 和 `mcp:tools` scope。
|
||||
依据企查查 OAuth 接入文档 V1.4,同一 grant 的 Access Token、Refresh Token 和
|
||||
client_id 可用于上述五个精确地址。连接前校验目标的 Protected Resource Metadata,
|
||||
五个 MCP 地址共用同一份 Bearer API Key。连接前校验目标的 Protected Resource Metadata,
|
||||
确认 resource 与 issuer 匹配;连接器不会向列表以外的 MCP 地址发送凭证。
|
||||
|
||||
## 共享授权与生命周期
|
||||
|
||||
- 一条连接器记录保存一份加密授权。每张卡片独立管理自己的账户凭证。
|
||||
- Octop 内部 HTTP MCP 网关聚合五类工具,保留输入 Schema,并加上类别前缀,避免工具重名。
|
||||
- 一条连接器记录保存一份加密 API Key。每张卡片独立管理自己的账户凭证。
|
||||
- Octop 以 `mcp_mode=internal` 托管内部 HTTP MCP:对话按 HTTP 加载 `/api/internal/mcp/qcc/…`,
|
||||
再聚合五类工具并加上类别前缀。不是进程内 gateway 适配器。
|
||||
一张卡片即可选用所有五类工具,无需创建五个独立连接器。
|
||||
- 每次上游请求读取最新凭证;临近到期时刷新,401 时最多刷新并重试一次。
|
||||
同一 Octop 进程、同一仓库实例内的并发请求共享刷新锁;轮换后的 Token 写回加密存储。
|
||||
此锁不覆盖多个 Octop 进程同时使用同一授权的部署。
|
||||
- 重启后从数据库恢复凭证和内部访问令牌,无需依赖原登录弹窗。授权失效时仍需重新授权。
|
||||
- 删除卡片时先撤销最新 Refresh Token,再删除本地凭证,五类服务一起解绑。
|
||||
撤销失败会保留凭证并报错,供用户重试。已发出的请求可能完成,且撤销 Refresh Token
|
||||
不等同于服务端立即吊销所有已签发 Access Token。
|
||||
- 探测会报告每类服务的结果;任一失败即不报告整体成功。运行时工具发现也要求五类服务均可访问。
|
||||
|
||||
## 本地授权与浏览器兼容性
|
||||
|
||||
本地部署的回调地址为 `http://127.0.0.1:<实际端口>/api/connectors/oauth/callback`
|
||||
(地址取决于访问 Octop 时使用的主机和端口)。保持 Octop 服务运行。
|
||||
若桌面端无法打开授权弹窗,可在独立 Chrome 浏览器中打开同一 Octop 本地页面授权。
|
||||
允许该 Octop 站点弹窗;如企查查页面请求访问本机应用或本地网络,按需允许该站点。
|
||||
|
||||
企查查成功页可能保留在弹窗中,由隐藏 iframe 触发本地回调。
|
||||
应以 Octop 的授权状态、连接探测和实际工具调用确认完成,不能仅以企查查成功页判断。
|
||||
Web/云端部署需要将实际 HTTPS 回调地址与企查查确认并加入白名单。
|
||||
- 删除卡片只清除本地凭证,不会远程吊销 API Key。
|
||||
- 探测会报告每类服务的结果;任一类别成功即整体可用。运行时工具发现同样只聚合成功类别。
|
||||
|
||||
## 验证范围
|
||||
|
||||
2026-09-23,贡献者在 Octop v1.0.1 的自定义 MCP 中,使用独立 Chrome、空 Headers
|
||||
完成 OAuth 授权,发现 company 的 16 个工具,并通过 `get_change_records` 返回实际企业变更数据。
|
||||
此记录验证了 company 服务与通用 OAuth 流程,不代表五类真实服务均已验证。
|
||||
2026-09-23,贡献者在 Octop 中使用合成凭证和 HTTP MockTransport,通过真实 MCP SDK
|
||||
验证五类服务的初始化、分页工具发现及工具调用;另外覆盖部分类别失败、内部接口鉴权
|
||||
及跨用户删除权限。测试不访问真实账户。
|
||||
|
||||
本变更的自动化测试使用合成凭证和 HTTP MockTransport,通过真实 MCP SDK 验证五类服务的
|
||||
初始化、分页工具发现及工具调用;另外覆盖并发刷新、重复 401、数据库重新打开后的凭证恢复、
|
||||
刷新与解绑竞争、撤销失败重试、内部接口鉴权及跨用户解绑权限。测试不访问真实账户。
|
||||
|
||||
五类服务的真实业务调用、真实 Token 到期刷新、完整应用重启和真实撤销仍需授权账户验收。
|
||||
桌面内嵌浏览器的弹窗及回调兼容不属于此次变更。工具数量由服务端决定,不作为固定契约。
|
||||
五类服务的真实业务调用仍需授权账户验收。工具数量由服务端决定,不作为固定契约。
|
||||
|
||||
官方入口:<https://agent.qcc.com/>。
|
||||
连接器图标来自该站点公开的 `/favicon-qcc.png`,用于标识企查查服务。
|
||||
|
||||
@@ -863,11 +863,7 @@ async def delete_instance(
|
||||
user: Any = Depends(current_user),
|
||||
server: Any = Depends(get_server),
|
||||
) -> None:
|
||||
"""Disconnect and delete stored credentials for a connector instance.
|
||||
|
||||
QCC revokes the shared refresh token for all five resources first. If remote
|
||||
revocation fails, credentials are retained so disconnection can be retried.
|
||||
"""
|
||||
"""Disconnect and delete stored credentials for a connector instance."""
|
||||
custom_target = _resolve_custom_target(instance_id, user=user, server=server)
|
||||
if custom_target is not None:
|
||||
custom_user_id, synthetic_name = custom_target
|
||||
@@ -903,13 +899,7 @@ async def delete_instance(
|
||||
cli_creds = {"instance_id": instance_id}
|
||||
else:
|
||||
cli_creds = {**cli_creds, "instance_id": instance_id}
|
||||
if inst.kind == "qcc":
|
||||
try:
|
||||
await _connector_service(server).disconnect_qcc(instance_id)
|
||||
except ValueError as exc:
|
||||
raise OctopError(ErrorCode.CONNECTOR_INVALID_CREDENTIALS, str(exc)) from exc
|
||||
else:
|
||||
repo.delete(instance_id)
|
||||
repo.delete(instance_id)
|
||||
if cli_creds is not None:
|
||||
cleanup_creds_cli_dirs(inst.kind, cli_creds)
|
||||
_schedule_connector_reload(server, user_id, all_users=inst.shared)
|
||||
|
||||
@@ -11,7 +11,9 @@ from octop.config import OctopConfig
|
||||
from octop.infra.connectors.catalog import (
|
||||
ConnectorCatalogEntry,
|
||||
get_catalog_entry,
|
||||
is_inprocess_gateway,
|
||||
is_mcp_oauth_remote,
|
||||
uses_internal_http_mcp,
|
||||
)
|
||||
from octop.infra.connectors.custom_mcp import validate_mcp_http_url
|
||||
from octop.infra.connectors.mail_servers import resolve_mail_servers
|
||||
@@ -82,8 +84,6 @@ def build_http_mcp_spec(
|
||||
creds: dict[str, Any],
|
||||
config: OctopConfig,
|
||||
) -> dict[str, Any]:
|
||||
if entry.kind == "qcc":
|
||||
return _build_gateway_spec(entry, instance_id, creds, config)
|
||||
if entry.mcp_mode == "remote":
|
||||
return _build_remote_spec(entry, creds)
|
||||
return _build_gateway_spec(entry, instance_id, creds, config)
|
||||
@@ -247,7 +247,7 @@ def validate_create_credentials(
|
||||
out["oauth_client_secret"] = str(credentials["oauth_client_secret"])
|
||||
if credentials.get("openid"):
|
||||
out["openid"] = str(credentials["openid"])
|
||||
if entry.mcp_mode == "gateway":
|
||||
if is_inprocess_gateway(entry) or uses_internal_http_mcp(entry):
|
||||
out["internal_token"] = new_internal_token()
|
||||
return out
|
||||
|
||||
@@ -495,7 +495,7 @@ def build_mcp_server_configs_for_user(
|
||||
)
|
||||
for inst, entry, creds in _iter_active_connectors(svc, connector_repo, user_id):
|
||||
try:
|
||||
if entry.mcp_mode == "gateway":
|
||||
if is_inprocess_gateway(entry):
|
||||
# Name-only placeholder: harness skips specs without ``transport``;
|
||||
# tools are injected in-process in AgentManager._post_start_agent.
|
||||
configs[inst.mcp_server_name] = {}
|
||||
@@ -545,18 +545,18 @@ def build_mcp_server_configs_for_user(
|
||||
|
||||
|
||||
def gateway_mcp_server_names(*, connector_repo: Any, user_id: int) -> set[str]:
|
||||
"""MCP server names of *user_id*'s active gateway-mode connector instances.
|
||||
"""MCP server names of in-process gateway connectors for *user_id*.
|
||||
|
||||
Gateway connectors carry no HTTP transport: their tools are built in-process
|
||||
from stored credentials, so callers can attach them to a live agent instead
|
||||
of rebuilding it.
|
||||
``mcp_mode=gateway`` has no HTTP transport: tools are injected from Python
|
||||
adapters. ``internal`` aggregators (e.g. QCC) are not included — harness
|
||||
loads those via ``/api/internal/mcp``.
|
||||
"""
|
||||
names: set[str] = set()
|
||||
for inst in connector_repo.list_visible(user_id):
|
||||
if inst.status != "active":
|
||||
continue
|
||||
entry = get_catalog_entry(inst.kind)
|
||||
if entry is not None and entry.mcp_mode == "gateway":
|
||||
if entry is not None and is_inprocess_gateway(entry):
|
||||
names.add(inst.mcp_server_name)
|
||||
return names
|
||||
|
||||
@@ -588,7 +588,7 @@ def inject_missing_gateway_tools(
|
||||
)
|
||||
return
|
||||
for inst, entry, creds in _iter_active_connectors(svc, connector_repo, user_id):
|
||||
if entry.mcp_mode != "gateway":
|
||||
if not is_inprocess_gateway(entry):
|
||||
continue
|
||||
if inst.mcp_server_name not in wanted:
|
||||
continue
|
||||
@@ -608,7 +608,7 @@ def inject_missing_gateway_tools(
|
||||
for inst in connector_repo.list_visible(user_id)
|
||||
if inst.status == "active"
|
||||
and (entry := get_catalog_entry(inst.kind)) is not None
|
||||
and entry.mcp_mode == "gateway"
|
||||
and is_inprocess_gateway(entry)
|
||||
]
|
||||
http_loaded = [
|
||||
n for n in gateway_names if any(str(t).startswith(f"{n}_") for t in tool_set)
|
||||
|
||||
@@ -18,6 +18,7 @@ AuthKind = Literal[
|
||||
|
||||
CredentialFieldType = Literal["text", "password", "url", "tags"]
|
||||
RemoteTransport = Literal["raw_http", "streamable_http", "sse"]
|
||||
McpMode = Literal["remote", "gateway", "internal"]
|
||||
ConnectorCategory = Literal[
|
||||
"office",
|
||||
"knowledge",
|
||||
@@ -50,7 +51,10 @@ class ConnectorCatalogEntry:
|
||||
icon: str
|
||||
color: str
|
||||
phase: Literal["available", "coming_soon"]
|
||||
mcp_mode: Literal["remote", "gateway"]
|
||||
# remote: harness talks to the vendor URL.
|
||||
# gateway: in-process Python adapter; harness config is a name-only placeholder.
|
||||
# internal: Octop-hosted HTTP MCP at /api/internal/mcp; harness loads via HTTP.
|
||||
mcp_mode: McpMode
|
||||
category: ConnectorCategory
|
||||
quick_auth_url: str | None = None
|
||||
login_url: str | None = None
|
||||
@@ -72,6 +76,16 @@ class ConnectorCatalogEntry:
|
||||
remote_transport: RemoteTransport = "raw_http"
|
||||
|
||||
|
||||
def is_inprocess_gateway(entry: ConnectorCatalogEntry) -> bool:
|
||||
"""Harness injects Python adapter tools; config is a name-only placeholder."""
|
||||
return entry.mcp_mode == "gateway"
|
||||
|
||||
|
||||
def uses_internal_http_mcp(entry: ConnectorCatalogEntry) -> bool:
|
||||
"""Harness loads Octop-hosted HTTP MCP at ``/api/internal/mcp``."""
|
||||
return entry.mcp_mode == "internal"
|
||||
|
||||
|
||||
def is_mcp_oauth_remote(entry: ConnectorCatalogEntry) -> bool:
|
||||
"""True when this catalog entry is a dynamic-OAuth remote MCP connector."""
|
||||
return (
|
||||
@@ -363,21 +377,18 @@ _CATALOG: tuple[ConnectorCatalogEntry, ...] = (
|
||||
ConnectorCatalogEntry(
|
||||
kind="qcc",
|
||||
name="企查查",
|
||||
description="一次授权接入企查查五类 MCP:企业数据、风险数据、知识产权、经营信息与董监高信息",
|
||||
auth_kind="oauth2",
|
||||
description="用一份 API Key 接入企查查五类 MCP:企业数据、风险数据、知识产权、经营信息与董监高信息",
|
||||
auth_kind="api_key",
|
||||
doc_url="https://agent.qcc.com/",
|
||||
icon="qcc",
|
||||
color="#008CFF",
|
||||
phase="available",
|
||||
mcp_mode="remote",
|
||||
mcp_mode="internal",
|
||||
category="professional",
|
||||
guide_url="https://agent.qcc.com/",
|
||||
auth_hint="点击「一键授权」登录企查查并授权,无需手动复制 API Key;查询范围与额度以企查查账户权限为准",
|
||||
oauth_issuer="https://agent.qcc.com",
|
||||
mcp_url="https://agent.qcc.com/mcp/company/stream",
|
||||
oauth_resource="https://agent.qcc.com/mcp/company/stream",
|
||||
oauth_scopes="mcp:tools",
|
||||
remote_transport="streamable_http",
|
||||
quick_auth_url="https://agent.qcc.com/",
|
||||
guide_url="https://agent.qcc.com/guide",
|
||||
manual_url="https://agent.qcc.com/",
|
||||
auth_hint="登录企查查智能体平台,在个人中心复制 API Key 并粘贴到下方;查询范围以账户权限为准,未开通的类别会跳过",
|
||||
),
|
||||
ConnectorCatalogEntry(
|
||||
kind="tencent-ardot",
|
||||
@@ -441,7 +452,7 @@ _CATALOG: tuple[ConnectorCatalogEntry, ...] = (
|
||||
mcp_mode="remote",
|
||||
category="knowledge",
|
||||
guide_url="https://help.openalex.org/access/connector/",
|
||||
auth_hint="点击「一键授权」登录 OpenAlex;查询将使用你自己的 API Key 与每日预算。",
|
||||
auth_hint="点击「一键授权」登录 OpenAlex(桌面端请用系统浏览器);查询将使用你自己的 API Key 与每日预算。",
|
||||
oauth_issuer="https://mcp.openalex.org",
|
||||
mcp_url="https://mcp.openalex.org/mcp",
|
||||
oauth_resource="https://mcp.openalex.org/mcp",
|
||||
|
||||
@@ -7,7 +7,7 @@ from typing import Any
|
||||
|
||||
from harness_agent.mcp import mcp_args_model, sanitize_llm_tool_name
|
||||
|
||||
from octop.infra.connectors.catalog import ConnectorCatalogEntry
|
||||
from octop.infra.connectors.catalog import ConnectorCatalogEntry, is_inprocess_gateway
|
||||
from octop.infra.connectors.gateway.protocol import handle_mcp_request
|
||||
from octop.infra.connectors.gateway.registry import mcp_tools_for_kind
|
||||
|
||||
@@ -23,7 +23,7 @@ def build_gateway_langchain_tools(
|
||||
from langchain_core.tools import StructuredTool
|
||||
|
||||
del instance_id
|
||||
if not isinstance(entry, ConnectorCatalogEntry) or entry.mcp_mode != "gateway":
|
||||
if not isinstance(entry, ConnectorCatalogEntry) or not is_inprocess_gateway(entry):
|
||||
return []
|
||||
out: list[Any] = []
|
||||
|
||||
|
||||
@@ -374,6 +374,11 @@ async def probe_connector(
|
||||
instance_id: str,
|
||||
config: OctopConfig,
|
||||
) -> dict[str, Any]:
|
||||
if entry.kind == "qcc":
|
||||
from octop.infra.connectors.qcc import bearer_token, probe
|
||||
|
||||
return await probe(bearer_token(cred_payload))
|
||||
|
||||
if entry.mcp_mode == "gateway":
|
||||
try:
|
||||
await asyncio.to_thread(probe_gateway_credentials, entry.kind, cred_payload)
|
||||
@@ -404,11 +409,6 @@ async def probe_connector(
|
||||
return {"ok": False, "error": "请填写 API Key"}
|
||||
return await probe_youdao_note(api_key)
|
||||
|
||||
if entry.kind == "qcc":
|
||||
from octop.infra.connectors.qcc import probe
|
||||
|
||||
return await probe(str(cred_payload.get("access_token") or ""))
|
||||
|
||||
spec = build_http_mcp_spec(
|
||||
entry=entry,
|
||||
instance_id=instance_id,
|
||||
|
||||
@@ -1,17 +1,14 @@
|
||||
"""One QCC grant, five fixed MCP resources, one namespaced tool surface."""
|
||||
"""One QCC API key, five fixed MCP resources, one namespaced tool surface."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import time
|
||||
from typing import Any
|
||||
|
||||
import httpx
|
||||
from mcp import ClientSession
|
||||
from mcp.client.streamable_http import streamable_http_client
|
||||
|
||||
from octop.infra.connectors.oauth.mcp import (
|
||||
_ensure_mcp_oauth_url,
|
||||
fetch_authorization_metadata,
|
||||
)
|
||||
from octop.infra.utils.ssrf_guard import safe_request
|
||||
|
||||
ISSUER = "https://agent.qcc.com"
|
||||
@@ -20,22 +17,39 @@ RESOURCES = {
|
||||
for name in ("company", "risk", "ipr", "operation", "executive")
|
||||
}
|
||||
|
||||
# When a resource actually exposes these names, hide the rest so the model is
|
||||
# not handed 150+ near-duplicate tools. Unknown catalogs (and unit doubles)
|
||||
# fall through to the full page.
|
||||
PREFERRED_TOOL_SUFFIXES: frozenset[str] = frozenset(
|
||||
{
|
||||
"get_company_profile",
|
||||
"get_change_records",
|
||||
"get_business_exception",
|
||||
"get_patent_info",
|
||||
"get_administrative_license",
|
||||
"get_executive_positions",
|
||||
}
|
||||
)
|
||||
|
||||
def unauthorized(exc: BaseException) -> bool:
|
||||
if isinstance(exc, BaseExceptionGroup):
|
||||
return any(unauthorized(child) for child in exc.exceptions)
|
||||
return isinstance(exc, httpx.HTTPStatusError) and exc.response.status_code == 401
|
||||
_METADATA_TTL_SEC = 600.0
|
||||
_metadata_ok_until: dict[str, float] = {}
|
||||
|
||||
|
||||
async def request_resource(
|
||||
resource: str, token: str, method: str, params: dict[str, Any]
|
||||
) -> dict[str, Any]:
|
||||
"""Use the MCP SDK for session initialization, pagination and tool calls.
|
||||
def bearer_token(creds: dict[str, Any]) -> str:
|
||||
return str(
|
||||
creds.get("api_key") or creds.get("token") or creds.get("access_token") or ""
|
||||
).strip()
|
||||
|
||||
No caller-supplied URLs and no credential-bearing redirects are permitted.
|
||||
Sessions are short-lived so a rotated grant is used on the next call.
|
||||
"""
|
||||
|
||||
def clear_metadata_cache() -> None:
|
||||
_metadata_ok_until.clear()
|
||||
|
||||
|
||||
async def _assert_resource_metadata(resource: str) -> str:
|
||||
url = RESOURCES[resource]
|
||||
now = time.monotonic()
|
||||
if _metadata_ok_until.get(resource, 0.0) > now:
|
||||
return url
|
||||
metadata_response = await safe_request(
|
||||
"GET",
|
||||
f"{ISSUER}/mcp/.well-known/oauth-protected-resource/{resource}/stream",
|
||||
@@ -45,6 +59,20 @@ async def request_resource(
|
||||
metadata = metadata_response.json()
|
||||
if metadata.get("resource") != url or ISSUER not in metadata.get("authorization_servers", []):
|
||||
raise ValueError("QCC resource metadata mismatch")
|
||||
_metadata_ok_until[resource] = now + _METADATA_TTL_SEC
|
||||
return url
|
||||
|
||||
|
||||
async def request_resource(
|
||||
resource: str, token: str, method: str, params: dict[str, Any]
|
||||
) -> dict[str, Any]:
|
||||
"""Use the MCP SDK for session initialization, pagination and tool calls.
|
||||
|
||||
No caller-supplied URLs and no credential-bearing redirects are permitted.
|
||||
Sessions are short-lived so the stored API key is read on the next call.
|
||||
Protected-resource metadata is cached so list/call do not re-handshake.
|
||||
"""
|
||||
url = await _assert_resource_metadata(resource)
|
||||
async with (
|
||||
httpx.AsyncClient(
|
||||
headers={"Authorization": f"Bearer {token}"},
|
||||
@@ -76,8 +104,23 @@ async def request_resource(
|
||||
seen.add(cursor)
|
||||
|
||||
|
||||
def _tool_suffix(name: str) -> str:
|
||||
_, sep, rest = name.partition("__")
|
||||
return rest if sep else name
|
||||
|
||||
|
||||
def select_exposed_tools(tools: list[dict[str, Any]]) -> list[dict[str, Any]]:
|
||||
matched = [
|
||||
tool
|
||||
for tool in tools
|
||||
if _tool_suffix(str(tool.get("name") or "")) in PREFERRED_TOOL_SUFFIXES
|
||||
]
|
||||
return matched or tools
|
||||
|
||||
|
||||
def namespace_tools(resource: str, result: dict[str, Any]) -> list[dict[str, Any]]:
|
||||
return [{**tool, "name": f"{resource}__{tool['name']}"} for tool in result.get("tools", [])]
|
||||
named = [{**tool, "name": f"{resource}__{tool['name']}"} for tool in result.get("tools", [])]
|
||||
return select_exposed_tools(named)
|
||||
|
||||
|
||||
async def probe(token: str) -> dict[str, Any]:
|
||||
@@ -94,32 +137,9 @@ async def probe(token: str) -> dict[str, Any]:
|
||||
servers[resource] = {"ok": False}
|
||||
failed = [name for name, result in servers.items() if not result["ok"]]
|
||||
return {
|
||||
"ok": not failed,
|
||||
"ok": bool(tools),
|
||||
"tools": tools,
|
||||
"tool_count": len(tools),
|
||||
"servers": servers,
|
||||
**({"error": f"QCC MCP probe failed: {', '.join(failed)}"} if failed else {}),
|
||||
}
|
||||
|
||||
|
||||
async def revoke(creds: dict[str, Any]) -> None:
|
||||
token = str(creds.get("refresh_token") or "")
|
||||
if not token:
|
||||
return
|
||||
metadata = await fetch_authorization_metadata(ISSUER)
|
||||
endpoint = await _ensure_mcp_oauth_url(
|
||||
str(metadata.get("revocation_endpoint") or ""),
|
||||
issuer=ISSUER,
|
||||
field="revocation_endpoint",
|
||||
)
|
||||
response = await safe_request(
|
||||
"POST",
|
||||
endpoint,
|
||||
data={
|
||||
"client_id": str(creds.get("oauth_client_id") or ""),
|
||||
"token": token,
|
||||
"token_type_hint": "refresh_token",
|
||||
},
|
||||
timeout=20.0,
|
||||
)
|
||||
response.raise_for_status()
|
||||
|
||||
@@ -7,12 +7,11 @@ import logging
|
||||
import time
|
||||
from pathlib import Path
|
||||
from typing import Any
|
||||
from weakref import WeakKeyDictionary
|
||||
|
||||
from octop.config import OctopConfig
|
||||
from octop.infra.connectors import qcc
|
||||
from octop.infra.connectors.builder import build_http_mcp_spec, mcp_server_name, new_internal_token
|
||||
from octop.infra.connectors.catalog import get_catalog_entry
|
||||
from octop.infra.connectors.catalog import get_catalog_entry, uses_internal_http_mcp
|
||||
from octop.infra.connectors.crypto import decrypt_credentials, encrypt_credentials
|
||||
from octop.infra.connectors.custom_mcp import (
|
||||
CUSTOM_MCP_DISPLAY_NAME,
|
||||
@@ -51,10 +50,6 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
_OAUTH_REFRESH_SKEW_SEC = 120
|
||||
|
||||
# Octop runs in one process. Services constructed by different routes share
|
||||
# the application's repository, and therefore the same per-grant lock.
|
||||
_QCC_LOCKS: WeakKeyDictionary[ConnectorRepo, dict[str, asyncio.Lock]] = WeakKeyDictionary()
|
||||
|
||||
|
||||
class ConnectorNameTakenError(ValueError):
|
||||
"""Raised when a connector display name is already owned by the user."""
|
||||
@@ -108,11 +103,7 @@ class ConnectorService:
|
||||
row = self._repo.get(instance_id)
|
||||
if row is None or not row.credential_blob:
|
||||
return {}
|
||||
creds = decrypt_credentials(self._secret_repo, row.credential_blob)
|
||||
if row.kind == "qcc" and not creds.get("internal_token"):
|
||||
creds["internal_token"] = new_internal_token()
|
||||
self.encrypt_and_store(instance_id=instance_id, payload=creds)
|
||||
return creds
|
||||
return decrypt_credentials(self._secret_repo, row.credential_blob)
|
||||
|
||||
def encrypt_and_store(
|
||||
self,
|
||||
@@ -122,17 +113,19 @@ class ConnectorService:
|
||||
) -> None:
|
||||
stored = dict(payload)
|
||||
row = self._repo.get(instance_id)
|
||||
if row is not None and row.kind == "qcc":
|
||||
existing = (
|
||||
decrypt_credentials(self._secret_repo, row.credential_blob)
|
||||
if row.credential_blob
|
||||
else {}
|
||||
)
|
||||
stored["internal_token"] = (
|
||||
existing.get("internal_token")
|
||||
or stored.get("internal_token")
|
||||
or new_internal_token()
|
||||
)
|
||||
if row is not None:
|
||||
entry = get_catalog_entry(row.kind)
|
||||
if entry is not None and uses_internal_http_mcp(entry):
|
||||
existing = (
|
||||
decrypt_credentials(self._secret_repo, row.credential_blob)
|
||||
if row.credential_blob
|
||||
else {}
|
||||
)
|
||||
stored["internal_token"] = (
|
||||
existing.get("internal_token")
|
||||
or stored.get("internal_token")
|
||||
or new_internal_token()
|
||||
)
|
||||
stored["instance_id"] = instance_id
|
||||
expires_at = stored.get("expires_at")
|
||||
exp = int(expires_at) if expires_at is not None else None
|
||||
@@ -145,7 +138,10 @@ class ConnectorService:
|
||||
kind: str,
|
||||
) -> dict[str, Any]:
|
||||
if kind == "qcc":
|
||||
return await self._fresh_qcc(instance_id)
|
||||
row = self._repo.get(instance_id)
|
||||
if row is None or row.kind != "qcc" or row.status != "active":
|
||||
return {}
|
||||
return self.decrypt(instance_id)
|
||||
creds = self.decrypt(instance_id)
|
||||
entry = get_catalog_entry(kind)
|
||||
if entry is None or entry.auth_kind != "oauth2":
|
||||
@@ -168,51 +164,11 @@ class ConnectorService:
|
||||
self.encrypt_and_store(instance_id=instance_id, payload=creds)
|
||||
return creds
|
||||
|
||||
def _qcc_lock(self, instance_id: str) -> asyncio.Lock:
|
||||
locks = _QCC_LOCKS.setdefault(self._repo, {})
|
||||
return locks.setdefault(instance_id, asyncio.Lock())
|
||||
|
||||
async def _fresh_qcc(
|
||||
self, instance_id: str, *, rejected_token: str | None = None
|
||||
) -> dict[str, Any]:
|
||||
async with self._qcc_lock(instance_id):
|
||||
row = self._repo.get(instance_id)
|
||||
if row is None or row.kind != "qcc" or row.status != "active":
|
||||
return {}
|
||||
creds = self.decrypt(instance_id)
|
||||
token = str(creds.get("access_token") or "")
|
||||
# Another request may already have rotated the token rejected by a 401.
|
||||
force = rejected_token is not None and rejected_token == token
|
||||
expires = int(creds.get("expires_at") or 0)
|
||||
if not force and expires > int(time.time()) + _OAUTH_REFRESH_SKEW_SEC:
|
||||
return creds
|
||||
if not creds.get("refresh_token"):
|
||||
if force or (expires and expires <= int(time.time())):
|
||||
raise ValueError("QCC authorization expired; please authorize again")
|
||||
return creds
|
||||
if not creds.get("oauth_client_id"):
|
||||
raise ValueError("QCC client registration missing; please authorize again")
|
||||
refreshed = await refresh_oauth_credentials(
|
||||
kind="qcc", creds=creds, settings_repo=self._settings_repo
|
||||
)
|
||||
creds.update(refreshed)
|
||||
self.encrypt_and_store(instance_id=instance_id, payload=creds)
|
||||
return creds
|
||||
|
||||
async def _qcc_request(
|
||||
self, instance_id: str, resource: str, method: str, params: dict[str, Any]
|
||||
) -> dict[str, Any]:
|
||||
creds = await self._fresh_qcc(instance_id)
|
||||
token = str(creds.get("access_token") or "")
|
||||
if not token:
|
||||
raise ValueError("QCC connector is disconnected")
|
||||
try:
|
||||
return await qcc.request_resource(resource, token, method, params)
|
||||
except Exception as exc:
|
||||
if not qcc.unauthorized(exc):
|
||||
raise
|
||||
creds = await self._fresh_qcc(instance_id, rejected_token=token)
|
||||
token = str(creds.get("access_token") or "")
|
||||
creds = await self.ensure_fresh_credentials(instance_id, "qcc")
|
||||
token = qcc.bearer_token(creds)
|
||||
if not token:
|
||||
raise ValueError("QCC connector is disconnected")
|
||||
return await qcc.request_resource(resource, token, method, params)
|
||||
@@ -225,11 +181,21 @@ class ConnectorService:
|
||||
return handle_mcp_request(kind="qcc", creds={}, body=body)
|
||||
try:
|
||||
if method == "tools/list":
|
||||
listed = await asyncio.gather(
|
||||
*[
|
||||
self._qcc_request(instance_id, resource, "tools/list", {})
|
||||
for resource in qcc.RESOURCES
|
||||
],
|
||||
return_exceptions=True,
|
||||
)
|
||||
tools: list[dict[str, Any]] = []
|
||||
for resource in qcc.RESOURCES:
|
||||
result = await self._qcc_request(instance_id, resource, "tools/list", {})
|
||||
tools.extend(qcc.namespace_tools(resource, result))
|
||||
result = {"tools": tools}
|
||||
for resource, item in zip(qcc.RESOURCES, listed, strict=True):
|
||||
if isinstance(item, BaseException):
|
||||
continue
|
||||
tools.extend(qcc.namespace_tools(resource, item))
|
||||
if not tools:
|
||||
raise ValueError("QCC MCP request failed; check API Key or try again")
|
||||
result: dict[str, Any] = {"tools": tools}
|
||||
else:
|
||||
params = dict(body.get("params") or {})
|
||||
resource, sep, name = str(params.get("name") or "").partition("__")
|
||||
@@ -246,24 +212,10 @@ class ConnectorService:
|
||||
"id": body.get("id"),
|
||||
"error": {
|
||||
"code": -32603,
|
||||
"message": "QCC MCP request failed; check connection or authorize again",
|
||||
"message": "QCC MCP request failed; check API Key or try again",
|
||||
},
|
||||
}
|
||||
|
||||
async def disconnect_qcc(self, instance_id: str) -> None:
|
||||
async with self._qcc_lock(instance_id):
|
||||
row = self._repo.get(instance_id)
|
||||
if row is None:
|
||||
return
|
||||
if row.kind != "qcc":
|
||||
raise ValueError("Not a QCC connector")
|
||||
try:
|
||||
await qcc.revoke(self.decrypt(instance_id))
|
||||
except Exception as exc:
|
||||
# Keep the encrypted grant so the user can retry remote revocation.
|
||||
raise ValueError("QCC revocation failed; retry disconnect") from exc
|
||||
self._repo.delete(instance_id)
|
||||
|
||||
def reserved_builtin_mcp_names(self, user_id: int) -> set[str]:
|
||||
names: set[str] = set()
|
||||
for inst in self._repo.list_by_user(user_id):
|
||||
|
||||
@@ -66,9 +66,10 @@ async def test_catalog(env):
|
||||
assert weiyun["category"] == "office"
|
||||
assert weiyun.get("quick_auth_url") == "https://www.weiyun.com/act/openclaw"
|
||||
qcc = next(e for e in r.json() if e["kind"] == "qcc")
|
||||
assert qcc["auth_kind"] == "oauth2"
|
||||
assert qcc["oauth_mode"] == "dynamic"
|
||||
assert qcc["oauth_ready"] is True
|
||||
assert qcc["auth_kind"] == "api_key"
|
||||
assert qcc["mcp_mode"] == "internal"
|
||||
assert qcc["oauth_mode"] is None
|
||||
assert qcc["oauth_ready"] is False
|
||||
assert qcc["category"] == "professional"
|
||||
openalex = next(e for e in r.json() if e["kind"] == "openalex")
|
||||
assert openalex == {
|
||||
@@ -86,7 +87,7 @@ async def test_catalog(env):
|
||||
"login_url": None,
|
||||
"guide_url": "https://help.openalex.org/access/connector/",
|
||||
"manual_url": "https://help.openalex.org/access/connector/",
|
||||
"auth_hint": "点击「一键授权」登录 OpenAlex;查询将使用你自己的 API Key 与每日预算。",
|
||||
"auth_hint": "点击「一键授权」登录 OpenAlex(桌面端请用系统浏览器);查询将使用你自己的 API Key 与每日预算。",
|
||||
"oauth_mode": "dynamic",
|
||||
"oauth_ready": True,
|
||||
"credential_fields": [],
|
||||
@@ -701,8 +702,6 @@ async def test_custom_mcp_oauth_start_unified(env):
|
||||
|
||||
|
||||
async def test_qcc_gateway_auth_five_resources_and_disconnect(env, monkeypatch):
|
||||
import time
|
||||
|
||||
from octop.api.routers.internal_mcp import _service
|
||||
from octop.infra.connectors import qcc
|
||||
|
||||
@@ -713,12 +712,7 @@ async def test_qcc_gateway_auth_five_resources_and_disconnect(env, monkeypatch):
|
||||
json={
|
||||
"kind": "qcc",
|
||||
"display_name": "QCC",
|
||||
"credentials": {
|
||||
"access_token": "synthetic",
|
||||
"refresh_token": "synthetic-refresh",
|
||||
"oauth_client_id": "client",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
},
|
||||
"credentials": {"api_key": "synthetic"},
|
||||
},
|
||||
)
|
||||
assert created.status_code == 201
|
||||
@@ -741,25 +735,12 @@ async def test_qcc_gateway_auth_five_resources_and_disconnect(env, monkeypatch):
|
||||
other = await create_user(c, auth, username="qcc_reader")
|
||||
denied = await c.delete(f"/api/connector-instances/{instance_id}", headers=other)
|
||||
assert denied.status_code == 403
|
||||
revoke = AsyncMock()
|
||||
monkeypatch.setattr(qcc, "revoke", revoke)
|
||||
deleted = await c.delete(f"/api/connector-instances/{instance_id}", headers=auth)
|
||||
assert deleted.status_code == 204
|
||||
revoke.assert_awaited_once()
|
||||
gone = await c.post(path, params={"token": token}, json={"id": 2, "method": "tools/list"})
|
||||
assert gone.status_code == 404
|
||||
|
||||
|
||||
async def test_qcc_disconnect_api_docs(tmp_octop_home):
|
||||
write_octop_config(tmp_octop_home, enable_api_docs=True)
|
||||
async with octop_client(tmp_octop_home) as (c, _):
|
||||
assert (await c.get("/api/docs")).status_code == 200
|
||||
schema = (await c.get("/api/openapi.json")).json()
|
||||
operation = schema["paths"]["/api/connector-instances/{instance_id}"]["delete"]
|
||||
assert operation["summary"] == "Delete connector"
|
||||
assert "five resources" in operation["description"]
|
||||
|
||||
|
||||
@pytest.mark.parametrize(
|
||||
("redirect_after", "expected_path"),
|
||||
[
|
||||
|
||||
@@ -1,33 +1,50 @@
|
||||
"""QCC catalog integration with the shared MCP OAuth and probe flows."""
|
||||
"""QCC catalog integration with API Key and the internal HTTP aggregator."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import hashlib
|
||||
from base64 import urlsafe_b64encode
|
||||
from pathlib import Path
|
||||
from unittest.mock import AsyncMock
|
||||
from urllib.parse import parse_qs, urlparse
|
||||
|
||||
import pytest
|
||||
|
||||
from octop.config import OctopConfig
|
||||
from octop.infra.connectors.builder import build_http_mcp_spec, validate_create_credentials
|
||||
from octop.infra.connectors.catalog import get_mcp_oauth_remote
|
||||
from octop.infra.connectors.oauth import registry
|
||||
from octop.infra.connectors.builder import (
|
||||
build_http_mcp_spec,
|
||||
build_mcp_server_configs_for_user,
|
||||
gateway_mcp_server_names,
|
||||
validate_create_credentials,
|
||||
)
|
||||
from octop.infra.connectors.catalog import (
|
||||
get_catalog_entry,
|
||||
get_mcp_oauth_remote,
|
||||
is_inprocess_gateway,
|
||||
uses_internal_http_mcp,
|
||||
)
|
||||
from octop.infra.connectors.probe import probe_connector
|
||||
|
||||
ISSUER = "https://agent.qcc.com"
|
||||
RESOURCE = f"{ISSUER}/mcp/company/stream"
|
||||
from octop.infra.connectors.service import ConnectorService
|
||||
from octop.infra.db.migrate import run_migrations
|
||||
from octop.infra.db.pool import SqlitePool
|
||||
from octop.infra.db.repos.connectors import ConnectorRepo
|
||||
from octop.infra.db.repos.secrets import SecretRepo
|
||||
from octop.infra.db.repos.settings import SettingsRepo
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_qcc_credentials_build_and_probe(monkeypatch: pytest.MonkeyPatch) -> None:
|
||||
entry = get_mcp_oauth_remote("qcc")
|
||||
entry = get_catalog_entry("qcc")
|
||||
assert entry is not None
|
||||
creds = validate_create_credentials("qcc", {"access_token": "test-qcc-token"})
|
||||
assert entry.mcp_mode == "internal"
|
||||
assert uses_internal_http_mcp(entry)
|
||||
assert not is_inprocess_gateway(entry)
|
||||
assert get_mcp_oauth_remote("qcc") is None
|
||||
creds = validate_create_credentials("qcc", {"api_key": "test-qcc-token"})
|
||||
assert creds["api_key"] == "test-qcc-token"
|
||||
assert creds["internal_token"]
|
||||
creds["internal_token"] = "local-gateway-token"
|
||||
spec = build_http_mcp_spec(
|
||||
entry=entry, instance_id="qcc-test", creds=creds, config=OctopConfig()
|
||||
)
|
||||
assert spec["transport"] == "http"
|
||||
assert "/api/internal/mcp/qcc/qcc-test?token=local-gateway-token" in spec["url"]
|
||||
assert "test-qcc-token" not in str(spec)
|
||||
probe = AsyncMock(return_value={"ok": True, "tools": []})
|
||||
@@ -37,68 +54,38 @@ async def test_qcc_credentials_build_and_probe(monkeypatch: pytest.MonkeyPatch)
|
||||
probe.assert_awaited_once_with("test-qcc-token")
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("metadata_scopes", [None, ["mcp:tools"]])
|
||||
async def test_qcc_catalog_authorization_parameters(
|
||||
monkeypatch: pytest.MonkeyPatch, metadata_scopes: list[str] | None
|
||||
) -> None:
|
||||
metadata = {
|
||||
"authorization_endpoint": f"{ISSUER}/oauth/authorize",
|
||||
"token_endpoint": f"{ISSUER}/oauth/token",
|
||||
"registration_endpoint": f"{ISSUER}/oauth/register",
|
||||
"token_endpoint_auth_methods_supported": ["none"],
|
||||
"scopes_supported": metadata_scopes,
|
||||
}
|
||||
fetch = AsyncMock(return_value=metadata)
|
||||
register = AsyncMock(return_value={"client_id": "test-public-client"})
|
||||
monkeypatch.setattr(registry, "fetch_authorization_metadata", fetch)
|
||||
monkeypatch.setattr(registry, "register_dynamic_client", register)
|
||||
callback = "http://127.0.0.1:8088/api/connectors/oauth/callback"
|
||||
url, verifier, ctx = await registry.start_oauth_for_target(
|
||||
target={"type": "catalog", "kind": "qcc"},
|
||||
redirect_uri=callback,
|
||||
state="test-state",
|
||||
settings_repo=None,
|
||||
)
|
||||
fetch.assert_awaited_once_with(ISSUER)
|
||||
assert register.await_args.kwargs["redirect_uri"] == callback
|
||||
query = parse_qs(urlparse(url).query)
|
||||
assert url.startswith(f"{ISSUER}/oauth/authorize?")
|
||||
assert query["client_id"] == ["test-public-client"]
|
||||
assert query["redirect_uri"] == [callback]
|
||||
assert query["scope"] == ["mcp:tools"]
|
||||
assert query["resource"] == [RESOURCE]
|
||||
assert query["state"] == ["test-state"]
|
||||
assert query["code_challenge_method"] == ["S256"]
|
||||
challenge = urlsafe_b64encode(hashlib.sha256(verifier.encode()).digest()).rstrip(b"=")
|
||||
assert query["code_challenge"] == [challenge.decode()]
|
||||
assert ctx["issuer"] == ISSUER
|
||||
assert ctx["resource"] == RESOURCE
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_qcc_refresh_uses_public_client_and_company_resource(
|
||||
monkeypatch: pytest.MonkeyPatch,
|
||||
) -> None:
|
||||
metadata = {"token_endpoint": f"{ISSUER}/oauth/token"}
|
||||
fetch = AsyncMock(return_value=metadata)
|
||||
refresh = AsyncMock(
|
||||
return_value={"access_token": "new-access", "refresh_token": "rotated-refresh"}
|
||||
)
|
||||
monkeypatch.setattr(registry, "fetch_authorization_metadata", fetch)
|
||||
monkeypatch.setattr(registry, "refresh_access_token", refresh)
|
||||
result = await registry.refresh_oauth_credentials(
|
||||
def test_qcc_harness_config_uses_http_not_inprocess_gateway(tmp_path: Path) -> None:
|
||||
pool = SqlitePool(tmp_path / "octop.db")
|
||||
run_migrations(pool)
|
||||
with pool.transaction() as conn:
|
||||
conn.execute(
|
||||
"INSERT INTO users(id,username,password_hash,role,created_at) VALUES (1,'qcc','x','user',1)"
|
||||
)
|
||||
repo = ConnectorRepo(pool)
|
||||
repo.create(
|
||||
instance_id="qcc1",
|
||||
user_id=1,
|
||||
kind="qcc",
|
||||
creds={"oauth_client_id": "test-public-client", "refresh_token": "old-refresh"},
|
||||
settings_repo=None,
|
||||
display_name="QCC",
|
||||
mcp_server_name="qcc__qcc1",
|
||||
)
|
||||
fetch.assert_awaited_once_with(ISSUER)
|
||||
refresh.assert_awaited_once_with(
|
||||
metadata,
|
||||
issuer=ISSUER,
|
||||
client_id="test-public-client",
|
||||
client_secret=None,
|
||||
refresh_token="old-refresh",
|
||||
resource=RESOURCE,
|
||||
svc = ConnectorService(
|
||||
repo=repo,
|
||||
secret_repo=SecretRepo(pool),
|
||||
settings_repo=SettingsRepo(pool),
|
||||
config=OctopConfig(),
|
||||
)
|
||||
assert result["refresh_token"] == "rotated-refresh"
|
||||
svc.encrypt_and_store(instance_id="qcc1", payload={"api_key": "k", "internal_token": "tok"})
|
||||
configs = build_mcp_server_configs_for_user(
|
||||
svc=svc,
|
||||
connector_repo=repo,
|
||||
user_id=1,
|
||||
agent_id="agent",
|
||||
agent_user_id=1,
|
||||
config=OctopConfig(),
|
||||
log=False,
|
||||
)
|
||||
spec = configs["qcc__qcc1"]
|
||||
assert spec.get("transport") == "http"
|
||||
assert "/api/internal/mcp/qcc/qcc1?token=tok" in spec["url"]
|
||||
assert "qcc__qcc1" not in gateway_mcp_server_names(connector_repo=repo, user_id=1)
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
"""Exercise the five real SDK transports with synthetic MCP HTTP responses."""
|
||||
"""Exercise the five real SDK transports with a stored API Key."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import time
|
||||
from pathlib import Path
|
||||
from unittest.mock import AsyncMock
|
||||
|
||||
@@ -43,15 +41,7 @@ def grant(tmp_path: Path):
|
||||
instance_id="grant", user_id=1, kind="qcc", display_name="QCC", mcp_server_name="qcc__grant"
|
||||
)
|
||||
svc = service(pool, repo)
|
||||
svc.encrypt_and_store(
|
||||
instance_id="grant",
|
||||
payload={
|
||||
"access_token": "old-access",
|
||||
"refresh_token": "old-refresh",
|
||||
"oauth_client_id": "public-client",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
},
|
||||
)
|
||||
svc.encrypt_and_store(instance_id="grant", payload={"api_key": "qcc-api-key"})
|
||||
return pool, repo, svc
|
||||
|
||||
|
||||
@@ -102,6 +92,7 @@ def mcp_http(monkeypatch: pytest.MonkeyPatch):
|
||||
json={"resource": qcc.RESOURCES[resource], "authorization_servers": [qcc.ISSUER]},
|
||||
)
|
||||
|
||||
qcc.clear_metadata_cache()
|
||||
monkeypatch.setattr(qcc, "safe_request", metadata)
|
||||
monkeypatch.setattr(qcc.httpx, "AsyncClient", factory)
|
||||
return calls
|
||||
@@ -125,67 +116,29 @@ async def test_five_servers_list_paginate_and_call(grant, mcp_http):
|
||||
)
|
||||
assert result["result"]["content"][0]["text"] == f"/mcp/{resource}/stream"
|
||||
assert {url for url, _, _ in mcp_http} == set(qcc.RESOURCES.values())
|
||||
assert {auth for _, auth, _ in mcp_http} == {"Bearer old-access"}
|
||||
assert {auth for _, auth, _ in mcp_http} == {"Bearer qcc-api-key"}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_concurrent_refresh_across_service_instances(grant, monkeypatch):
|
||||
pool, repo, svc = grant
|
||||
creds = svc.decrypt("grant")
|
||||
creds["expires_at"] = 1
|
||||
svc.encrypt_and_store(instance_id="grant", payload=creds)
|
||||
entered = asyncio.Event()
|
||||
release = asyncio.Event()
|
||||
|
||||
async def refresh(**kwargs):
|
||||
assert kwargs["creds"]["refresh_token"] == "old-refresh"
|
||||
entered.set()
|
||||
await release.wait()
|
||||
return {
|
||||
"access_token": "new-access",
|
||||
"refresh_token": "new-refresh",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
}
|
||||
|
||||
mocked = AsyncMock(side_effect=refresh)
|
||||
monkeypatch.setattr("octop.infra.connectors.service.refresh_oauth_credentials", mocked)
|
||||
tasks = [
|
||||
asyncio.create_task(service(pool, repo).ensure_fresh_credentials("grant", "qcc"))
|
||||
for _ in range(10)
|
||||
]
|
||||
await entered.wait()
|
||||
release.set()
|
||||
results = await asyncio.gather(*tasks)
|
||||
assert mocked.await_count == 1
|
||||
assert all(r["access_token"] == "new-access" for r in results)
|
||||
assert svc.decrypt("grant")["refresh_token"] == "new-refresh"
|
||||
assert b"new-refresh" not in repo.get("grant").credential_blob
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_restart_restores_rotated_grant(grant, mcp_http):
|
||||
async def test_restart_restores_api_key(grant, mcp_http):
|
||||
pool, _, svc = grant
|
||||
creds = svc.decrypt("grant")
|
||||
stable_token = creds["internal_token"]
|
||||
creds.update(access_token="rotated-access", refresh_token="rotated-refresh")
|
||||
creds["api_key"] = "rotated-key"
|
||||
svc.encrypt_and_store(instance_id="grant", payload=creds)
|
||||
# Reopen the on-disk DB with new repositories and no in-memory OAuth state.
|
||||
restarted = service(SqlitePool(pool.path))
|
||||
assert restarted.decrypt("grant")["internal_token"] == stable_token
|
||||
assert restarted.decrypt("grant")["refresh_token"] == "rotated-refresh"
|
||||
assert restarted.decrypt("grant")["api_key"] == "rotated-key"
|
||||
result = await restarted.handle_qcc_request("grant", {"id": 1, "method": "tools/list"})
|
||||
assert len(result["result"]["tools"]) == 10
|
||||
assert {auth for _, auth, _ in mcp_http} == {"Bearer rotated-access"}
|
||||
assert {auth for _, auth, _ in mcp_http} == {"Bearer rotated-key"}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_disconnect_revokes_latest_refresh_and_removes_all_servers(grant, monkeypatch):
|
||||
async def test_delete_removes_card_and_gateway_access(grant, monkeypatch):
|
||||
_, repo, svc = grant
|
||||
token = svc.decrypt("grant")["internal_token"]
|
||||
revoke = AsyncMock()
|
||||
monkeypatch.setattr(qcc, "revoke", revoke)
|
||||
await svc.disconnect_qcc("grant")
|
||||
assert revoke.await_args.args[0]["refresh_token"] == "old-refresh"
|
||||
repo.delete("grant")
|
||||
assert repo.get("grant") is None
|
||||
assert svc.verify_internal_token("grant", token) is None
|
||||
assert await svc.mcp_configs_for_user(1) == {}
|
||||
@@ -196,41 +149,6 @@ async def test_disconnect_revokes_latest_refresh_and_removes_all_servers(grant,
|
||||
request.assert_not_awaited()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_failed_revoke_keeps_grant_for_retry(grant, monkeypatch):
|
||||
_, repo, svc = grant
|
||||
monkeypatch.setattr(qcc, "revoke", AsyncMock(side_effect=ValueError("offline")))
|
||||
with pytest.raises(ValueError, match="retry disconnect"):
|
||||
await svc.disconnect_qcc("grant")
|
||||
assert repo.get("grant") is not None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_401_refreshes_once_and_retries(grant, monkeypatch):
|
||||
_, _, svc = grant
|
||||
exc = httpx.HTTPStatusError(
|
||||
"unauthorized",
|
||||
request=httpx.Request("POST", qcc.RESOURCES["risk"]),
|
||||
response=httpx.Response(401),
|
||||
)
|
||||
request = AsyncMock(side_effect=[exc, {"content": [], "isError": False}])
|
||||
refresh = AsyncMock(
|
||||
return_value={
|
||||
"access_token": "new-access",
|
||||
"refresh_token": "new-refresh",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
}
|
||||
)
|
||||
monkeypatch.setattr(qcc, "request_resource", request)
|
||||
monkeypatch.setattr("octop.infra.connectors.service.refresh_oauth_credentials", refresh)
|
||||
result = await svc.handle_qcc_request(
|
||||
"grant", {"id": 1, "method": "tools/call", "params": {"name": "risk__lookup"}}
|
||||
)
|
||||
assert "result" in result
|
||||
assert [c.args[1] for c in request.await_args_list] == ["old-access", "new-access"]
|
||||
assert refresh.await_count == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_unknown_resource_never_receives_token(grant, monkeypatch):
|
||||
_, _, svc = grant
|
||||
@@ -244,55 +162,21 @@ async def test_unknown_resource_never_receives_token(grant, monkeypatch):
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_disconnect_waits_for_rotation_without_resurrecting_grant(grant, monkeypatch):
|
||||
pool, repo, svc = grant
|
||||
creds = svc.decrypt("grant")
|
||||
creds["expires_at"] = 1
|
||||
svc.encrypt_and_store(instance_id="grant", payload=creds)
|
||||
entered, release = asyncio.Event(), asyncio.Event()
|
||||
|
||||
async def refresh(**kwargs):
|
||||
entered.set()
|
||||
await release.wait()
|
||||
return {
|
||||
"access_token": "new",
|
||||
"refresh_token": "rotated",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
}
|
||||
|
||||
monkeypatch.setattr("octop.infra.connectors.service.refresh_oauth_credentials", refresh)
|
||||
revoke = AsyncMock()
|
||||
monkeypatch.setattr(qcc, "revoke", revoke)
|
||||
rotating = asyncio.create_task(svc.ensure_fresh_credentials("grant", "qcc"))
|
||||
await entered.wait()
|
||||
deleting = asyncio.create_task(service(pool, repo).disconnect_qcc("grant"))
|
||||
await asyncio.sleep(0)
|
||||
revoke.assert_not_awaited()
|
||||
release.set()
|
||||
await asyncio.gather(rotating, deleting)
|
||||
assert revoke.await_args.args[0]["refresh_token"] == "rotated"
|
||||
assert repo.get("grant") is None
|
||||
assert await svc.ensure_fresh_credentials("grant", "qcc") == {}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_repeated_401_stops_after_one_retry(grant, monkeypatch):
|
||||
async def test_invalid_key_does_not_retry(grant, monkeypatch):
|
||||
_, _, svc = grant
|
||||
exc = httpx.HTTPStatusError(
|
||||
"secret-in-error",
|
||||
request=httpx.Request("POST", qcc.RESOURCES["risk"]),
|
||||
response=httpx.Response(401),
|
||||
request = AsyncMock(
|
||||
side_effect=httpx.HTTPStatusError(
|
||||
"secret-in-error",
|
||||
request=httpx.Request("POST", qcc.RESOURCES["risk"]),
|
||||
response=httpx.Response(401),
|
||||
)
|
||||
)
|
||||
request = AsyncMock(side_effect=exc)
|
||||
refresh = AsyncMock(return_value={"access_token": "new", "expires_at": int(time.time()) + 3600})
|
||||
monkeypatch.setattr(qcc, "request_resource", request)
|
||||
monkeypatch.setattr("octop.infra.connectors.service.refresh_oauth_credentials", refresh)
|
||||
result = await svc.handle_qcc_request(
|
||||
"grant", {"id": 1, "method": "tools/call", "params": {"name": "risk__lookup"}}
|
||||
)
|
||||
assert "error" in result and "secret-in-error" not in json.dumps(result)
|
||||
assert request.await_count == 2
|
||||
assert refresh.await_count == 1
|
||||
assert request.await_count == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@@ -304,39 +188,67 @@ async def test_probe_reports_partial_failure(monkeypatch):
|
||||
|
||||
monkeypatch.setattr(qcc, "request_resource", request)
|
||||
result = await qcc.probe("synthetic-token")
|
||||
assert result["ok"] is False
|
||||
assert result["ok"] is True
|
||||
assert result["tool_count"] == 4
|
||||
assert result["servers"]["risk"] == {"ok": False}
|
||||
assert len(result["servers"]) == 5
|
||||
assert "secret" not in json.dumps(result)
|
||||
|
||||
|
||||
def test_select_exposed_tools_prefers_known_suffixes():
|
||||
preferred = [
|
||||
{"name": "company__get_company_profile"},
|
||||
{"name": "company__lookup"},
|
||||
{"name": "company__get_change_records"},
|
||||
]
|
||||
assert [tool["name"] for tool in qcc.select_exposed_tools(preferred)] == [
|
||||
"company__get_company_profile",
|
||||
"company__get_change_records",
|
||||
]
|
||||
unknown = [{"name": "company__lookup"}, {"name": "company__detail"}]
|
||||
assert qcc.select_exposed_tools(unknown) == unknown
|
||||
|
||||
|
||||
def test_bearer_token_prefers_api_key():
|
||||
assert qcc.bearer_token({"api_key": "k", "access_token": "legacy"}) == "k"
|
||||
assert qcc.bearer_token({"token": "t"}) == "t"
|
||||
assert qcc.bearer_token({}) == ""
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_revoke_discovers_validates_and_posts_refresh_token(monkeypatch):
|
||||
endpoint = qcc.ISSUER + "/oauth/revoke"
|
||||
metadata = AsyncMock(return_value={"revocation_endpoint": endpoint})
|
||||
validate = AsyncMock(return_value=endpoint)
|
||||
request = AsyncMock(return_value=httpx.Response(200, request=httpx.Request("POST", endpoint)))
|
||||
monkeypatch.setattr(qcc, "fetch_authorization_metadata", metadata)
|
||||
monkeypatch.setattr(qcc, "_ensure_mcp_oauth_url", validate)
|
||||
monkeypatch.setattr(qcc, "safe_request", request)
|
||||
await qcc.revoke({"refresh_token": "latest-refresh", "oauth_client_id": "client"})
|
||||
metadata.assert_awaited_once_with(qcc.ISSUER)
|
||||
validate.assert_awaited_once_with(endpoint, issuer=qcc.ISSUER, field="revocation_endpoint")
|
||||
assert request.await_args.kwargs["data"] == {
|
||||
"client_id": "client",
|
||||
"token": "latest-refresh",
|
||||
"token_type_hint": "refresh_token",
|
||||
}
|
||||
validate.side_effect = ValueError("untrusted endpoint")
|
||||
request.reset_mock()
|
||||
with pytest.raises(ValueError):
|
||||
await qcc.revoke({"refresh_token": "latest-refresh"})
|
||||
request.assert_not_awaited()
|
||||
async def test_list_keeps_tools_when_one_resource_fails(grant, monkeypatch):
|
||||
_, _, svc = grant
|
||||
|
||||
async def request(resource, *_args, **_kwargs):
|
||||
if resource == "risk":
|
||||
raise ValueError("denied")
|
||||
return {"tools": [{"name": "lookup"}]}
|
||||
|
||||
monkeypatch.setattr(qcc, "request_resource", request)
|
||||
listed = await svc.handle_qcc_request("grant", {"id": 1, "method": "tools/list"})
|
||||
names = [tool["name"] for tool in listed["result"]["tools"]]
|
||||
assert "risk__lookup" not in names
|
||||
assert names == [f"{resource}__lookup" for resource in qcc.RESOURCES if resource != "risk"]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_resource_metadata_is_cached(mcp_http, monkeypatch):
|
||||
fetches: list[str] = []
|
||||
original = qcc.safe_request
|
||||
|
||||
async def counted(method, url, **kwargs):
|
||||
fetches.append(url)
|
||||
return await original(method, url, **kwargs)
|
||||
|
||||
monkeypatch.setattr(qcc, "safe_request", counted)
|
||||
await qcc.request_resource("risk", "old-access", "tools/list", {})
|
||||
await qcc.request_resource("risk", "old-access", "tools/list", {})
|
||||
assert len(fetches) == 1
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_metadata_mismatch_blocks_bearer_transport(monkeypatch):
|
||||
qcc.clear_metadata_cache()
|
||||
response = httpx.Response(
|
||||
200,
|
||||
request=httpx.Request("GET", qcc.ISSUER),
|
||||
@@ -348,43 +260,3 @@ async def test_metadata_mismatch_blocks_bearer_transport(monkeypatch):
|
||||
with pytest.raises(ValueError, match="metadata mismatch"):
|
||||
await qcc.request_resource("risk", "secret", "tools/list", {})
|
||||
transport.assert_not_called()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_concurrent_401_requests_share_rotated_token(grant, monkeypatch):
|
||||
pool, repo, _ = grant
|
||||
arrived = 0
|
||||
all_arrived = asyncio.Event()
|
||||
|
||||
async def request(resource, token, method, params):
|
||||
nonlocal arrived
|
||||
if token == "old-access":
|
||||
arrived += 1
|
||||
if arrived == 5:
|
||||
all_arrived.set()
|
||||
await all_arrived.wait()
|
||||
raise httpx.HTTPStatusError(
|
||||
"expired",
|
||||
request=httpx.Request("POST", qcc.RESOURCES[resource]),
|
||||
response=httpx.Response(401),
|
||||
)
|
||||
assert token == "new-access"
|
||||
return {"content": []}
|
||||
|
||||
refresh = AsyncMock(
|
||||
return_value={
|
||||
"access_token": "new-access",
|
||||
"refresh_token": "new-refresh",
|
||||
"expires_at": int(time.time()) + 3600,
|
||||
}
|
||||
)
|
||||
monkeypatch.setattr(qcc, "request_resource", request)
|
||||
monkeypatch.setattr("octop.infra.connectors.service.refresh_oauth_credentials", refresh)
|
||||
results = await asyncio.gather(
|
||||
*(
|
||||
service(pool, repo)._qcc_request("grant", resource, "tools/call", {"name": "lookup"})
|
||||
for resource in qcc.RESOURCES
|
||||
)
|
||||
)
|
||||
assert all(result == {"content": []} for result in results)
|
||||
refresh.assert_awaited_once()
|
||||
|
||||
Reference in New Issue
Block a user