b4d04eb8fd
* feat: add session-scoped connector access for onboarding
* fix(connectors): availability is the config flag AND the portal entitlement — no free-tier leg
The port carried a third availability leg from hermes-magic: a stored guest
(free-tier) identity short-circuits the managed-tool entitlement check. That
leg reads hermes_cli.anon_auth, which does not exist on hermes-agent main, so
connectors_available() raised ImportError inside its fail-closed try and the
whole connector surface was silently dark on a plain upstream checkout.
On this tree availability is the two-leg AND the design started with:
tools.connectors.enabled AND managed_nous_tools_enabled(). The free-tier leg
is a hermes-magic concern and belongs in hermes-magic's own delta over this
branch, next to the identity it depends on. Its integration test goes with it.
* docs(tool-search): connectors section — remote tools through the bridge
The squashed port carried the code but not the user-facing docs. Restores the
Connectors section of the Tool Search page and the connector-gateway host /
CONNECTOR_GATEWAY_URL override on the Tool Gateway page, updated for the
manage_connections tool and the pure-connector batch rule.
* fix(tool-search): connector tools rank with local tools in one pass instead of taking leftover slots
dispatch_tool_search ran BM25 over the local catalog, filled `limit` slots,
then appended connector hits only into slots left empty. On a 300-tool
catalog no slot was ever empty, so with Gmail and Google Calendar connected
"send gmail email" returned five betterstack tools and zero connector tools.
The gateway's hits for a query now become catalog entries (connector name,
slug words, description as the search text) and join the local catalog for
that query's BM25 pass. One ranking, one rarest-token admission rule for both
sources, `limit` as the total per query. The merge loop and the separate
record builder for connector hits are gone; `_shared_tool_record` serves both
sources.
The gateway search timeout rises from 8 s to 30 s. One request with six
use_cases measured 7 s, so 8 s sat on the edge and cut real answers off; the
failure path is unchanged (local-only results, no error to the model).
Live, 311 local tools + gateway, before -> after:
"send gmail email": 5 betterstack tools -> gmail SEND_EMAIL, CREATE_EMAIL_DRAFT
"read google calendar events": 5 betterstack tools -> googlecalendar EVENTS_LIST_ALL_CALENDARS
"linear create issue", "betterstack incident": unchanged
Benchmark (25 labelled queries): connector recall 0.09 -> 0.82, precision@5
0.18 -> 0.59, false positives on absent intents 17 -> 2.
* refactor(tool-search): connector leg into tools/connector_search.py
tools/tool_search.py is a facade. The connector leg (gateway hits as catalog
entries for tool_search, remote schemas for tool_describe, the
connections_in_scope gate) was appended to it by the port. It now lives in
its own sibling, tools/connector_search.py, and the facade imports the three
entry points: connections_in_scope, connector_entries_by_group,
remote_schemas_for.
No behaviour change. The tool_describe remote block became
remote_schemas_for(names, current_tool_defs, connector_describe) with the
same inputs, the same silent-degradation contract and the same injection
seam the tests already use.
* fix(tool-search): at most 7 queries per call, the gateway's search limit
One tool_search call sends all its queries to the connector gateway as one
search request. The gateway answers 7 use_cases per request and returns
HTTP 502 for 8 or more (measured 2026-09-09, re-measured with one-word
use_cases: it is a count limit, not a size limit). With the client cap at
10, a model sending 8 to 10 queries lost every connector hit for that call
and saw local-only results with no error.
The shared constant splits: _MAX_QUERIES_PER_CALL = 7 for search,
_MAX_DESCRIBE_NAMES_PER_CALL = 10 for describe, which has no remote count
limit. Eight or more queries now get the existing "too many queries" retry
hint before any request is made. No chunking: one call, one request.
* fix(tool-search): the model is told that connectors__ names are manage_connections accounts
tool_search results carry names like connectors__gmail__CREATE_EMAIL_DRAFT and
manage_connections is the tool that checks and connects those accounts, but
nothing told the model the two are the same thing. A model that hit
CONNECTION_REQUIRED had to infer the fix on its own.
The tool_search description gains one sentence making the link, added at
assembly only when manage_connections is in the session's tools. Signed out
or with connectors off the tool is absent and the description is unchanged,
so it never names a tool the model cannot call. This follows the existing
rule for cross-tool references (tools/AGENTS.md): they are added dynamically
from the session's actual tool set, never hardcoded in a schema.
Tool defs are fixed for the life of a conversation, so the description is
byte-stable per conversation; this is a one-time prefix change.
Live, real get_tool_definitions() against a signed-in home: sentence present.
Same home with auth.json removed: manage_connections absent, sentence absent.
* fix(connectors): /stop halts a connector batch before the next remote call
dispatch_connector_batch runs every remote entry of a tool_call batch in
sequence. The executor only checks the interrupt flag between tools, and
the whole batch is one tool to it, so a /stop landing during entry 1 of
20 still sent the other 19 to the gateway.
The loop now reads tools.interrupt.is_interrupted before each dispatch.
Once set, it stops calling handle_function_call and fills every unstarted
slot with the loop's existing error-slot shape, code INTERRUPTED and the
message "Stopped by the user before this call was made.", so the result
envelope stays valid and the counts stay honest. Entries already
dispatched keep their real results.
Test: three connector calls where the fake client sets the interrupt on
the first execute. The client sees exactly one call and slots 2 and 3
carry INTERRUPTED. Red on the base branch, green with the fix.
* test(connections): schema assertions become dispatch contracts
test_schema_documents_wait_and_its_timeout froze description fragments
("REQUIRED", "can NOT disconnect", "Nous Portal"). A wording edit fails
it while a real regression (a disconnect that reaches the gateway) does
not. That is a snapshot of prose, not a behaviour contract.
Delete it. The requirement that wait needs connectors is already covered
by test_wait_requires_connectors. The user-only disconnect boundary is
now asserted as behaviour: action disconnect with a connector returns an
error and the fake client records no call. That replaces the earlier
de-authenticate test, which only checked that the word "dashboard"
appeared in the error text.
Test count in the file goes from 26 to 25.
* docs(tool-search): connector batches are one gateway request per entry
The user guide said a connector batch travels as one gateway request. It
does not: model_tools_connectors.dispatch_connector_batch re-enters core
dispatch per entry, and each entry becomes its own execute request in
bridge._run_remote (plus at most one literal-slug retry when the gateway
reports TOOL_NOT_FOUND under the conventional slug). The docstrings in
tools/tool_gateway/bridge.py and tools/tool_gateway/__init__.py still
described the abandoned V1 plan and claimed nothing outside the package
imports it.
Rewrite those sentences to match the code: one request per entry, in
input order, dispatched from model_tools_connectors.py, with the per-entry
approval and interrupt behaviour that motivated the split. The guide also
still showed the single-call shape tool_call(name, arguments); both
places now show the `calls: [{name, arguments}]` array the schema
advertises and note that a single local call is an array of one.
Docs only, no test.
* fix(tools): the between-turns refresh never rewrites the bridge tools
The per-turn MCP refresh folds a fresh tool snapshot into the live array
with preserve_prefix: order and membership stay, but a name present in both
takes the fresh schema. That is right for ordinary tools, whose schema is a
constant. tool_search is the one tool whose description is derived from the
session: the deferred-tool count, the embedded listing, and, on this branch,
whether manage_connections was present. A late MCP server or one failed
portal lookup (manage_connections' check_fn fails closed) changed those bytes
on the next turn, and every byte after tool_search in the cached prefix was
re-prefilled. The array also contradicted itself in that case: the flapping
manage_connections was carried forward while the description lost its hint.
The bridge entries now keep the bytes they were built with for the life of
the conversation. Nothing is lost: tool_search reads the live catalog at
dispatch, so tools that arrived late are still found; connector availability
is checked at dispatch too. The compaction-boundary rebuild (content_aware,
the one sanctioned cache break) still refreshes the description.
Consequence: connector exposure in the prompt is decided once, at agent
build, by whether the user was signed in then. That is the intended
contract.
* refactor(tool-search): normalize_tool_call_entries lives with the other argument validation
The port appended the tool_call argument parser to the tool_search facade.
The family already has tools/tool_search_validation.py for exactly this
work (schema validation of deferred call arguments), so the parser moves
there and the facade imports it. No behaviour change; the one test that
imported it now imports from the defining module.
* refactor(connectors): delete the unused batch dispatcher; _run_remote becomes run_remote
bridge.dispatch_calls and its helpers (_dispatch_calls_inner, _run_pre_dispatch,
_run_local, _error_slot, _maybe_parse_json) and the LocalDispatch / PreDispatch
seams had no production caller. Connector dispatch runs through
model_tools_connectors: dispatch_connector_batch re-enters handle_function_call
once per entry, so scope, hook, approval and middleware policy fire against each
composed name inside core dispatch, and dispatch_connector_call hands the single
planned entry to the bridge's transport function. Only tests called the batch
dispatcher, and they exercised policy seams that production never wires.
The transport function is the module's real entry point, so it drops the
underscore: _run_remote becomes run_remote, body unchanged. The module
docstring now describes the two legs that exist (availability with D32 silent
degradation, and run_remote) instead of the injected seams. Imports that only
the deleted code used are gone; merge.py is untouched because every export
still has a caller.
Tests that drove dispatch_calls are deleted where they covered the removed
seams (pre_dispatch blocks and rewrites, local_dispatch classification, mixed
batches). The literal-slug fallback, the per-entry transport failure, and the
hook rewrite reaching the gateway request body are re-targeted at
handle_function_call('tool_call', ...) with the fake client swapped in at
bridge._default_client_factory, the same seam test_connector_dispatch_policy
uses. Each re-targeted test fails when the retry is disabled in run_remote.
* fix(connectors): search keeps the twin a colliding name reaches, and says so
format_connector_name strips the toolkit prefix, so GMAIL_FETCH_PROFILE and a
literal FETCH_PROFILE on gmail both compose to connectors__gmail__FETCH_PROFILE.
describe and execute decode that name to the prefixed slug first, so the
literal twin is unreachable under it. If a vendor ever shipped both, search
could describe the literal under a name that runs the prefixed tool.
Search is the one place that sees both twins in one response. It now keeps
the twin the name reaches and drops the other with a WARNING that names both
slugs, whichever the gateway listed first. Short names stay; no marker, no
per-process map, no change to describe or execute. No such pair exists in the
live catalog today; the guard turns a silent alias into a logged one.
307 lines
12 KiB
Python
307 lines
12 KiB
Python
"""HTTP client for the connector routes on the managed tool gateway.
|
|
|
|
Constructed per dispatch — a portal access token expires within the hour, so
|
|
auth headers are read fresh on every call (the ``managed_gateway_auth_headers``
|
|
idiom). Sync ``requests`` on purpose: the bridge branch cannot reach the
|
|
registry's async bridge, so every call here runs on the calling thread.
|
|
|
|
Injectable seams (``transport`` / ``endpoint_resolver`` / ``header_provider``)
|
|
default to the real ones; tests inject fakes instead of patching modules.
|
|
|
|
Retry policy (D29): at most ONE retry, on transport failure or 5xx only,
|
|
reusing the SAME ``x-idempotency-key``. The key is a local variable in
|
|
:meth:`ConnectorClient.execute` — one dispatch is one call frame, nothing
|
|
outside it ever retries the same dispatch, so scope guarantees same-key-on-
|
|
retry with no store to clean up. 409 means the key was reused with a
|
|
different body — always a client bug, never retried.
|
|
|
|
Only the gateway's execute route supports idempotency keys, so only routes
|
|
that are safe to repeat are retried: search and schemas are read-only, and
|
|
execute dedupes on its key. The connections route is state-changing WITHOUT
|
|
dedup support (it starts or restarts an authorization flow), so it is never
|
|
retried automatically — a lost response surfaces as an error the caller can
|
|
deliberately re-ask.
|
|
|
|
Wire models stay inside this module: callers receive plain dicts shaped for
|
|
``merge.splice_remote_results``.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import logging
|
|
import uuid
|
|
from typing import Any, Callable, Optional, Protocol, Sequence
|
|
|
|
import requests
|
|
|
|
from tools.tool_gateway import wire
|
|
from tools.tool_gateway.errors import (
|
|
GatewayAuthError,
|
|
GatewayUnavailable,
|
|
ToolGatewayError,
|
|
parse_gateway_error,
|
|
)
|
|
from tools.tool_gateway.merge import PlannedCall
|
|
|
|
logger = logging.getLogger(__name__)
|
|
|
|
__all__ = ["ConnectorClient", "Transport"]
|
|
|
|
DEFAULT_TIMEOUT_SECONDS = 30.0
|
|
# One execute request carries up to MAX_CALLS_PER_DISPATCH remote tool runs;
|
|
# measured batch latency is seconds, not minutes, but give slow tools room.
|
|
EXECUTE_TIMEOUT_SECONDS = 60.0
|
|
# Search rides the availability path of EVERY tool_search once connectors are
|
|
# lit. A hung gateway degrades silently to local-only results, with no retry.
|
|
# Measured: one request with 6 use_cases takes about 7 s, so an 8 s budget sat
|
|
# on the edge and cut real answers off; 30 s tolerates a slow gateway and still
|
|
# bounds the wait. Schemas (tool_describe) is user-initiated; a short budget
|
|
# with one retry keeps its worst case at 2x this value.
|
|
SEARCH_TIMEOUT_SECONDS = 30.0
|
|
SCHEMAS_TIMEOUT_SECONDS = 10.0
|
|
|
|
_MAX_RETRIES = 1 # D29: at most one retry, same key.
|
|
|
|
|
|
class Transport(Protocol):
|
|
"""The slice of ``requests`` the client uses; tests inject a fake."""
|
|
|
|
def request(
|
|
self,
|
|
method: str,
|
|
url: str,
|
|
*,
|
|
headers: Optional[dict] = None,
|
|
json: Optional[dict] = None,
|
|
timeout: Optional[float] = None,
|
|
) -> Any: ...
|
|
|
|
|
|
def _default_transport() -> Transport:
|
|
return requests # module satisfies the protocol
|
|
|
|
|
|
def _default_endpoint_resolver() -> Optional[str]:
|
|
"""The connectors origin, or ``None`` when none resolves (scheme misconfig).
|
|
|
|
Asks the connectors host resolver directly. This used to go through
|
|
``managed_vendor_endpoints("connectors")``, which invented a vendor that
|
|
does not exist — and then discarded the ``base_url``/``upload_path`` it
|
|
built for it. Connector routes are their own deployment's own paths
|
|
(``v1/connectors/*``) on its own host, not a vendor passthrough and not the
|
|
media host.
|
|
"""
|
|
from tools.managed_gateway_auth import connector_gateway_origin
|
|
|
|
try:
|
|
return connector_gateway_origin() or None
|
|
except ValueError:
|
|
# Misconfigured TOOL_GATEWAY_SCHEME: there is no origin to call.
|
|
return None
|
|
|
|
|
|
def _default_header_provider(url: str) -> dict:
|
|
from tools.managed_gateway_auth import managed_gateway_auth_headers
|
|
|
|
return managed_gateway_auth_headers(url)
|
|
|
|
|
|
class ConnectorClient:
|
|
"""One dispatch's connection to the gateway's connector routes."""
|
|
|
|
def __init__(
|
|
self,
|
|
*,
|
|
transport: Optional[Transport] = None,
|
|
endpoint_resolver: Optional[Callable[[], Optional[str]]] = None,
|
|
header_provider: Optional[Callable[[str], dict]] = None,
|
|
) -> None:
|
|
self._transport = transport or _default_transport()
|
|
self._endpoint_resolver = endpoint_resolver or _default_endpoint_resolver
|
|
self._header_provider = header_provider or _default_header_provider
|
|
|
|
# -- routes ---------------------------------------------------------
|
|
|
|
def search(self, queries: Sequence[dict[str, Any]]) -> dict[str, Any]:
|
|
"""POST v1/connectors/search. Returns the response as a plain dict."""
|
|
body = wire.ConnectorSearchRequest(
|
|
queries=[wire.ConnectorSearchQuery(**q) for q in queries]
|
|
).model_dump(by_alias=True, exclude_none=True)
|
|
payload = self._post(
|
|
wire.CONNECTOR_SEARCH_PATH, body,
|
|
timeout=SEARCH_TIMEOUT_SECONDS, retries=0,
|
|
)
|
|
parsed = wire.ConnectorSearchResponse.model_validate(payload)
|
|
return parsed.model_dump()
|
|
|
|
def schemas(self, tools: Sequence[str]) -> dict[str, Any]:
|
|
"""POST v1/connectors/schemas."""
|
|
body = wire.ConnectorSchemasRequest(tools=list(tools)).model_dump(
|
|
by_alias=True
|
|
)
|
|
payload = self._post(
|
|
wire.CONNECTOR_SCHEMAS_PATH, body, timeout=SCHEMAS_TIMEOUT_SECONDS
|
|
)
|
|
return wire.ConnectorSchemasResponse.model_validate(payload).model_dump()
|
|
|
|
def connections(
|
|
self, connectors: Sequence[str], *, reinitiate: bool = False
|
|
) -> dict[str, Any]:
|
|
"""POST v1/connectors/connections. Never auto-retried: this route
|
|
starts/restarts authorization flows and the gateway offers no dedup
|
|
key for it — a blind retry could double-submit a restart."""
|
|
body = wire.ConnectorConnectionsRequest(
|
|
connectors=list(connectors), reinitiate=reinitiate
|
|
).model_dump(by_alias=True)
|
|
payload = self._post(wire.CONNECTOR_CONNECTIONS_PATH, body, retries=0)
|
|
return wire.ConnectorConnectionsResponse.model_validate(payload).model_dump()
|
|
|
|
def list_connectors(self) -> list[dict[str, Any]]:
|
|
"""GET v1/connectors, following pagination. Read-only.
|
|
|
|
Returns the raw item dicts (``{"connector", "enabled", "connected"}``
|
|
today; tolerant of additions). The page size cap is the gateway's.
|
|
"""
|
|
items: list[dict[str, Any]] = []
|
|
cursor: Optional[str] = None
|
|
for _ in range(20): # generous page bound; the catalog is small
|
|
path = f"{wire.CONNECTORS_PATH}?limit=50"
|
|
if cursor:
|
|
path += f"&cursor={cursor}"
|
|
payload = self._request("GET", path, None)
|
|
if not isinstance(payload, dict):
|
|
break
|
|
page = payload.get("items")
|
|
if isinstance(page, list):
|
|
items.extend(entry for entry in page if isinstance(entry, dict))
|
|
cursor = payload.get("nextCursor")
|
|
if not cursor:
|
|
break
|
|
return items
|
|
|
|
def execute(self, planned: Sequence[PlannedCall]) -> list[dict[str, Any]]:
|
|
"""POST v1/connectors/execute — ONE request for the whole slice.
|
|
|
|
Returns one dict per wire result, in wire order (slot ``i`` is the
|
|
response to request ``tools[i]``): ``{"data": ..., "error": None |
|
|
{code, message, connector, connect_url, hint}}``. Length mismatches
|
|
are the merge layer's problem, by design.
|
|
"""
|
|
body = wire.ConnectorExecuteRequest(
|
|
tools=[
|
|
wire.ConnectorExecuteCall(
|
|
connector=plan.connector, tool=plan.tool, arguments=plan.arguments
|
|
)
|
|
for plan in planned
|
|
]
|
|
).model_dump(by_alias=True)
|
|
# Local by design (no store): one dispatch = one call frame, and a
|
|
# transport retry below re-presents this same key by scope.
|
|
idempotency_key = str(uuid.uuid4())
|
|
payload = self._post(
|
|
wire.CONNECTOR_EXECUTE_PATH,
|
|
body,
|
|
timeout=EXECUTE_TIMEOUT_SECONDS,
|
|
idempotency_key=idempotency_key,
|
|
)
|
|
parsed = wire.ConnectorExecuteResponse.model_validate(payload)
|
|
return [_result_dict(result) for result in parsed.results]
|
|
|
|
# -- plumbing -------------------------------------------------------
|
|
|
|
def _post(
|
|
self,
|
|
path: str,
|
|
body: dict[str, Any],
|
|
*,
|
|
timeout: float = DEFAULT_TIMEOUT_SECONDS,
|
|
idempotency_key: Optional[str] = None,
|
|
retries: int = _MAX_RETRIES,
|
|
) -> Any:
|
|
return self._request(
|
|
"POST", path, body,
|
|
timeout=timeout, idempotency_key=idempotency_key, retries=retries,
|
|
)
|
|
|
|
def _request(
|
|
self,
|
|
method: str,
|
|
path: str,
|
|
body: Optional[dict[str, Any]],
|
|
*,
|
|
timeout: float = DEFAULT_TIMEOUT_SECONDS,
|
|
idempotency_key: Optional[str] = None,
|
|
retries: int = _MAX_RETRIES,
|
|
) -> Any:
|
|
origin = self._endpoint_resolver()
|
|
if not origin:
|
|
raise GatewayUnavailable(
|
|
"no tool gateway origin resolves", code="NO_ORIGIN"
|
|
)
|
|
url = f"{origin.rstrip('/')}/{path}"
|
|
|
|
last_error: Optional[ToolGatewayError] = None
|
|
for attempt in range(1 + retries):
|
|
headers = dict(self._header_provider(url))
|
|
if not headers:
|
|
# No usable portal token; an unauthenticated request would
|
|
# only 401 — fail fast with the same meaning.
|
|
raise GatewayAuthError(
|
|
"no portal access token available", code="NO_TOKEN", status=401
|
|
)
|
|
headers["Content-Type"] = "application/json"
|
|
if idempotency_key:
|
|
headers["x-idempotency-key"] = idempotency_key
|
|
|
|
try:
|
|
response = self._transport.request(
|
|
method, url, headers=headers, json=body, timeout=timeout
|
|
)
|
|
except Exception as exc:
|
|
last_error = ToolGatewayError(
|
|
f"transport failure: {exc}", code="TRANSPORT_ERROR", retryable=True
|
|
)
|
|
logger.debug(
|
|
"Connector %s attempt %d transport failure: %s", path, attempt + 1, exc
|
|
)
|
|
continue
|
|
|
|
status = int(getattr(response, "status_code", 0))
|
|
if 200 <= status < 300:
|
|
return response.json()
|
|
|
|
error = parse_gateway_error(status, _safe_json(response))
|
|
if error.retryable and attempt < retries:
|
|
last_error = error
|
|
logger.debug(
|
|
"Connector %s attempt %d got %d; retrying with same key",
|
|
path,
|
|
attempt + 1,
|
|
status,
|
|
)
|
|
continue
|
|
raise error
|
|
|
|
assert last_error is not None # loop ran at least once
|
|
raise last_error
|
|
|
|
|
|
def _result_dict(result: wire.ConnectorExecuteResult) -> dict[str, Any]:
|
|
error = None
|
|
if result.error is not None:
|
|
error = {"code": result.error.code, "message": result.error.message}
|
|
if result.error.connector:
|
|
error["connector"] = result.error.connector
|
|
if result.error.connect_url:
|
|
error["connect_url"] = result.error.connect_url
|
|
if result.error.hint:
|
|
error["hint"] = result.error.hint
|
|
return {"data": result.data, "error": error}
|
|
|
|
|
|
def _safe_json(response: Any) -> Any:
|
|
try:
|
|
return response.json()
|
|
except Exception:
|
|
return getattr(response, "text", None)
|