diff --git a/apps/desktop/src/components/assistant-ui/connector-tool.tsx b/apps/desktop/src/components/assistant-ui/connector-tool.tsx index 643833a363..31ed231ae1 100644 --- a/apps/desktop/src/components/assistant-ui/connector-tool.tsx +++ b/apps/desktop/src/components/assistant-ui/connector-tool.tsx @@ -264,6 +264,10 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) { const { t } = useI18n() const copy = t.connectors const [reissuing, setReissuing] = useState>(new Set()) + // Rows whose link the user opened from this card. `initiated` only means a link was minted: the + // connect-first handoff (D85) mints on the watcher's first pass, before anyone clicks, so the + // "Waiting for your browser…" cue belongs to a row the user actually opened. + const [opened, setOpened] = useState>(new Set()) const unresolved = request.targets.some(target => !CONNECTOR_CARD_PHASES[target.state].resolved) // A DOM handle for the focus handoff, never rendered state. const cardRef = useRef(null) @@ -274,6 +278,12 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) { // A refused re-mint is a click that changed nothing, so it gets a toast; the row stays as it was. const reissue = async (target: ConnectionTarget): Promise => { setReissuing(current => new Set(current).add(target.name)) + setOpened(current => { + const next = new Set(current) + next.delete(target.name) + + return next + }) try { await reissueConnectionTarget(owner, request, target.name) @@ -319,6 +329,7 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) { {request.targets.map(target => { const phase = CONNECTOR_CARD_PHASES[target.state] const busy = reissuing.has(target.name) + const mark = phase.mark === 'waiting' && !opened.has(target.name) ? 'idle' : phase.mark const action = phase.verb === 'none' @@ -331,6 +342,7 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) { onClick: () => { if (phase.verb === 'open' && target.connectUrl && window.hermesDesktop?.openExternal) { void window.hermesDesktop.openExternal(target.connectUrl) + setOpened(current => new Set(current).add(target.name)) } if (phase.verb === 'reissue') { @@ -347,10 +359,10 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) { name: target.name, title: connectorTitle(target.name) }} - cue={phase.mark === 'waiting' ? copy.waiting : undefined} + cue={mark === 'waiting' ? copy.waiting : undefined} key={target.name} - mark={phase.mark} - markLabel={MARK_LABEL[phase.mark](copy)} + mark={mark} + markLabel={MARK_LABEL[mark](copy)} /> ) })} diff --git a/apps/desktop/src/components/onboarding-chat/cards/build.tsx b/apps/desktop/src/components/onboarding-chat/cards/build.tsx index 1457200518..85862205f1 100644 --- a/apps/desktop/src/components/onboarding-chat/cards/build.tsx +++ b/apps/desktop/src/components/onboarding-chat/cards/build.tsx @@ -14,6 +14,7 @@ import { quarantineHandoffReceipt } from '@/app/contrib/handoff-receipt' import { resolveSessionOwner } from '@/app/session/hooks/use-session-actions/utils' import type { CardProps } from '@/components/onboarding-chat/cards/frame' import { Chip } from '@/components/onboarding-chat/chip' +import { readPersistedHandoff } from '@/components/onboarding-chat/persisted-handoff' import { $handoffError, $setupHandoff, @@ -123,6 +124,8 @@ export function HandoffCard({ attrs, locked }: CardProps) { const brief = (attrs.brief ?? '').trim().slice(0, 240) const plan = parseHandoffPlan(attrs.plan) const state = useStore($setupHandoff) + // `locked` follows this text part; a later part (a tool call, reasoning) settles it while the reply still runs. + const replyRunning = useAuiState(s => s.message.status?.type === 'running') const receipt = useMemo(() => { try { @@ -139,24 +142,33 @@ export function HandoffCard({ attrs, locked }: CardProps) { const completed = receipt.completed useEffect(() => { - if (!task || !brief || locked || !storedId || !runtimeId || $setupHandoff.get() || completed) { + if (!task || !brief || locked || replyRunning || !storedId || !runtimeId || $setupHandoff.get() || completed) { return } let cancelled = false void resolveSessionOwner(storedId) - .then(owner => { + .then(async owner => { assertSessionOwnerResolved(owner, { method: 'onboarding.handoff', sessionId: storedId }) + const connectionId = isSessionOwnerRoute(owner) ? owner.connectionId : null + + const profile = isSessionOwnerRoute(owner) + ? owner.profile + : owner || $setupSession.get()?.profile || $activeGatewayProfile.get() + + // An unreachable history keeps the rendered attrs: today's behaviour, never a stalled handoff. + const persisted = await readPersistedHandoff(connectionId, profile, runtimeId).catch(() => null) + const persistedTask = (persisted?.task ?? '').trim().slice(0, 60) + const persistedBrief = (persisted?.brief ?? '').trim().slice(0, 240) + if (!cancelled) { - requestSetupHandoff(task, brief, plan, { - storedId, - runtimeId, - connectionId: isSessionOwnerRoute(owner) ? owner.connectionId : null, - profile: isSessionOwnerRoute(owner) - ? owner.profile - : owner || $setupSession.get()?.profile || $activeGatewayProfile.get() - }) + requestSetupHandoff( + persistedTask || task, + persistedBrief || brief, + persisted ? parseHandoffPlan(persisted.plan) : plan, + { storedId, runtimeId, connectionId, profile } + ) } }) .catch(error => { @@ -169,7 +181,7 @@ export function HandoffCard({ attrs, locked }: CardProps) { return () => { cancelled = true } - }, [brief, locked, plan, task, storedId, runtimeId, completed]) + }, [brief, locked, plan, replyRunning, task, storedId, runtimeId, completed]) if (!task || !brief) { return null diff --git a/apps/desktop/src/components/onboarding-chat/persisted-handoff.ts b/apps/desktop/src/components/onboarding-chat/persisted-handoff.ts new file mode 100644 index 0000000000..565ef7926c --- /dev/null +++ b/apps/desktop/src/components/onboarding-chat/persisted-handoff.ts @@ -0,0 +1,47 @@ +import type { SessionHistoryResult } from '@hermes/shared' + +import { segmentTranscriptDirectives } from '@/lib/transcript-directives' +import { requestGatewayForAgent } from '@/store/gateway' + +/** + * The handoff directive as the backend persisted it for the welcome chat. The build session is named and seeded + * from these attrs, so they come from the stored reply, not the renderer's streamed copy of it: on Windows the + * streamed copy once repeated its own chunks ("Set up mySet up my games…") while the stored reply was intact, and + * the handoff carried the garbled task and brief into the new session. Null when no persisted reply holds one. + */ +export async function readPersistedHandoff( + connectionId: null | string, + profile: string, + runtimeId: string +): Promise>> { + const history = await requestGatewayForAgent>( + connectionId, + profile, + 'session.history', + { + session_id: runtimeId + } + ) + + const messages = Array.isArray(history?.messages) ? history.messages : [] + + for (let index = messages.length - 1; index >= 0; index -= 1) { + const message = messages[index] + + if (message?.role !== 'assistant' || message.text == null) { + continue + } + + for (const segment of segmentTranscriptDirectives(message.text) ?? []) { + if ( + segment.kind === 'directive' && + segment.directive.name === 'onboarding' && + segment.directive.attrs.step === 'handoff' + ) { + return segment.directive.attrs + } + } + } + + return null +} diff --git a/hermes_cli/plugins_activation.py b/hermes_cli/plugins_activation.py index 54e91c8d35..869827e353 100644 --- a/hermes_cli/plugins_activation.py +++ b/hermes_cli/plugins_activation.py @@ -13,11 +13,14 @@ config or tree change. from __future__ import annotations import logging +import threading from pathlib import Path from typing import Any, Dict, List, Optional logger = logging.getLogger(__name__) +_GO_LIVE_LOCK = threading.Lock() + # Hooks the gateway consults per inbound/outbound message: live as soon as the registry holds them. _GATEWAY_TRANSFORM_HOOKS = frozenset({ "transform_llm_output", "transform_tool_result", "transform_terminal_output", "pre_gateway_dispatch", @@ -126,18 +129,34 @@ def load_and_go_live(name: str) -> Optional[Dict[str, Any]]: MCP servers and hand them (and its skills) to this profile's open chats with a turn note. Returns the activation summary with ``live_now: {mcp_servers, skills}``; ``deferred`` then keeps only what waits for the next session (Python ``tools``, ``prompt``). None when the plugin did not load.""" + # A forced rediscovery unloads every plugin before it loads them again, which drops their server + # configs, skills and liveness declarations for the length of the pass. Two installs finishing + # together (one card, two rows) each go live; the second one's pass must not run while the first + # reads or connects, or the first plugin comes up with no tools. One go-live at a time. + with _GO_LIVE_LOCK: + return _go_live(name) + + +def _go_live(name: str) -> Optional[Dict[str, Any]]: + from hermes_cli.plugins_activation_live import connect_plugin_mcp, live_notice, plugin_skills try: - from hermes_cli.plugins import discover_plugins, get_plugin_manager - discover_plugins(force=True) - activation = find_activation(activation_summaries(get_plugin_manager()), name) + from hermes_cli.plugins import _join_background_discovery, get_plugin_manager + _join_background_discovery() + manager = get_plugin_manager() + # Other forced passes (a reload-plugins verb, the dashboard) do not take the go-live lock, so the + # reads share the discovery lock with the pass that produced them. + with manager._discovery_lock: + manager.discover_and_load(force=True) + activation = find_activation(activation_summaries(manager), name) + portable = manager.get_portable_mcp_servers() + skills = plugin_skills(activation["key"]) if activation else [] except Exception: logger.debug("in-process plugin reload after change to %r failed", name, exc_info=True) return None if activation is None: return None - from hermes_cli.plugins_activation_live import connect_plugin_mcp, live_notice, plugin_skills - servers = connect_plugin_mcp(activation) - activation["live_now"] = {"mcp_servers": servers, "skills": plugin_skills(activation["key"])} + servers = connect_plugin_mcp(activation, portable) + activation["live_now"] = {"mcp_servers": servers, "skills": skills} activation["deferred"] = {k: v for k, v in (activation.get("deferred") or {}).items() if k != "mcp_servers"} import sys server = sys.modules.get("tui_gateway.server") # loaded == this process hosts chats diff --git a/hermes_cli/plugins_activation_live.py b/hermes_cli/plugins_activation_live.py index feb6356aa8..89cd0b5490 100644 --- a/hermes_cli/plugins_activation_live.py +++ b/hermes_cli/plugins_activation_live.py @@ -17,8 +17,9 @@ from typing import Any, Dict, List, Optional logger = logging.getLogger(__name__) -def connect_plugin_mcp(activation: Dict[str, Any]) -> List[Dict[str, Any]]: +def connect_plugin_mcp(activation: Dict[str, Any], portable: Dict[str, Dict[str, Any]]) -> List[Dict[str, Any]]: """Connect every portable MCP server ``activation`` defers, under the caller's profile scope. + ``portable`` is the manager's portable server configs, read together with ``activation``. Returns ``[{name, connected, tools, error?}]``, one row per server; a server that fails to connect is reported, never retried. Never raises.""" @@ -26,11 +27,9 @@ def connect_plugin_mcp(activation: Dict[str, Any]) -> List[Dict[str, Any]]: if not names: return [] try: - from hermes_cli.plugins import get_plugin_manager from tools.mcp_tool_config import _filter_suspicious_mcp_servers, _load_mcp_config from tools.mcp_tool_discovery import register_mcp_servers configured = _load_mcp_config() - portable = get_plugin_manager().get_portable_mcp_servers() except Exception as exc: logger.warning("plugin MCP activation could not read server config: %s", exc) return [{"name": n, "connected": False, "tools": [], "error": str(exc)} for n in names] diff --git a/hermes_cli/plugins_cmd.py b/hermes_cli/plugins_cmd.py index cc6bf8a302..11b662da12 100644 --- a/hermes_cli/plugins_cmd.py +++ b/hermes_cli/plugins_cmd.py @@ -11,9 +11,11 @@ import shutil import subprocess import sys import tempfile +import threading import urllib.parse +from contextlib import contextmanager from pathlib import Path -from typing import Any, NoReturn, Optional +from typing import Any, Callable, NoReturn, Optional from hermes_constants import get_hermes_home from hermes_cli._subprocess_compat import noninteractive_git_env @@ -560,6 +562,38 @@ def _write_install_metadata(metadata: dict[str, dict[str, object]]) -> None: path, json.dumps(metadata, indent=2, sort_keys=True) + "\n", tmp_prefix=f"{path.name}.tmp-") +_INSTALL_METADATA_LOCK_HOLDER = threading.local() + + +@contextmanager +def _install_metadata_lock(): + """Serialize read-modify-write of the sidecar across threads and processes. Installs overlap (the + Desktop install card runs its rows a second apart); each held a snapshot read before its clone, so + the later write dropped the earlier plugin's record.""" + from hermes_cli.auth import _file_lock + + path = _install_metadata_path() + with _file_lock(path.with_name(f"{path.name}.lock"), _INSTALL_METADATA_LOCK_HOLDER, 10.0, + "Timed out waiting for the plugin install metadata lock"): + yield + + +def _update_install_record(name: str, update: Callable[[Optional[dict]], Optional[dict]]) -> None: + """Rewrite one plugin's record in the CURRENT sidecar, under the lock. *update* maps the current + record (None when absent) to the new one (None removes it); every other record is re-read here, + never carried over from a caller's earlier snapshot.""" + with _install_metadata_lock(): + metadata = _read_install_metadata() + record = update(metadata.get(name)) + if record is None: + if name not in metadata: + return + del metadata[name] + else: + metadata[name] = record + _write_install_metadata(metadata) + + def pinned_revision(name: str, metadata: Optional[dict] = None) -> Optional[str]: """Full SHA a ``--ref`` install of *name* is pinned to, else ``None``.""" entry = (metadata if metadata is not None else _read_install_metadata()).get(name) @@ -1137,16 +1171,14 @@ def _post_pull_housekeeping(target: Path, console) -> None: def _remove_plugin_core(target: Path) -> None: """Remove one plugin and its metadata without splitting their state.""" - metadata = _read_install_metadata() - if target.name not in metadata: + if target.name not in _read_install_metadata(): rmtree_readonly(target) return - updated = {k: v for k, v in metadata.items() if k != target.name} staging = Path(tempfile.mkdtemp(prefix=f".{target.name}.remove-", dir=target.parent)) backup = staging / "plugin" os.replace(target, backup) try: - _write_install_metadata(updated) + _update_install_record(target.name, lambda _current: None) except Exception: try: os.replace(backup, target) diff --git a/hermes_cli/plugins_cmd_catalog.py b/hermes_cli/plugins_cmd_catalog.py index fa88d55c67..d92d0ebecf 100644 --- a/hermes_cli/plugins_cmd_catalog.py +++ b/hermes_cli/plugins_cmd_catalog.py @@ -99,14 +99,18 @@ def _install_record(plugin_dir: Path) -> Optional[dict]: def _write_catalog_block(plugin_dir: Path, record: dict, block: dict) -> dict: """Migrate one trusted installer record to the nested catalog contract.""" - from hermes_cli.plugins_cmd import _read_install_metadata, _write_install_metadata - migrated = dict(record) - migrated["catalog"] = block - migrated.pop("catalog_name", None) - migrated.pop("catalog_tier", None) - metadata = _read_install_metadata() - metadata[plugin_dir.name] = migrated - _write_install_metadata(metadata) + from hermes_cli.plugins_cmd import _update_install_record + + def migrate(current: Optional[dict]) -> Optional[dict]: + if current is None: + return None + migrated = dict(current) + migrated["catalog"] = block + migrated.pop("catalog_name", None) + migrated.pop("catalog_tier", None) + return migrated + + _update_install_record(plugin_dir.name, migrate) return block diff --git a/plugin-catalog/nvidia-app.yaml b/plugin-catalog/nvidia-app.yaml index 6a1e3ff832..c4b22498dd 100644 --- a/plugin-catalog/nvidia-app.yaml +++ b/plugin-catalog/nvidia-app.yaml @@ -1,6 +1,6 @@ name: nvidia-app repo: https://github.com/NousResearch/hermes-nvidia -sha: 5023f3a4ea5e1b55679090c598b39f04cf4aa848 +sha: 65bb8d89263f920e8599432823d2d4e05cc113e3 subdir: nvidia-app description: "NVIDIA App control through its local MCP server: overlay capture and state, driver status and release notes, game optimization, and app launch. Requires NVIDIA App 11.0 or newer on Windows; diff --git a/plugin-catalog/nvidia-broadcast.yaml b/plugin-catalog/nvidia-broadcast.yaml index f30eb92831..ef3d0383af 100644 --- a/plugin-catalog/nvidia-broadcast.yaml +++ b/plugin-catalog/nvidia-broadcast.yaml @@ -1,6 +1,6 @@ name: nvidia-broadcast repo: https://github.com/NousResearch/hermes-nvidia -sha: 5023f3a4ea5e1b55679090c598b39f04cf4aa848 +sha: 65bb8d89263f920e8599432823d2d4e05cc113e3 subdir: nvidia-broadcast description: "NVIDIA Broadcast control through its local MCP gateway: camera, microphone and speaker effects, Studio Voice profiles, device selection, camera resolution, and effects on local media files. diff --git a/pm/publication.py b/pm/publication.py index b50639867a..a2b42314d9 100644 --- a/pm/publication.py +++ b/pm/publication.py @@ -9,6 +9,7 @@ import base64 import hashlib import io import json +import threading from pathlib import Path from pm.environments import dependency_home_root, install_state_dir, runtime_facts_path @@ -16,6 +17,21 @@ from hermes_cli.runtime_state import _atomic_bytes, _bytes, _digest from pm.workspace import enabled_plugin_dirs, _is_member_candidate +_METADATA_LOCK_HOLDER = threading.local() + + +def _metadata_records(data: bytes | None) -> dict: + if data is None: + return {} + try: + records = json.loads(data) + except (UnicodeDecodeError, json.JSONDecodeError) as exc: + raise ValueError("Plugin install metadata changed while preparing the update; retry.") from exc + if not isinstance(records, dict): + raise ValueError("Plugin install metadata changed while preparing the update; retry.") + return records + + def candidate_members(extra_dirs=(), **selection): selected = enabled_plugin_dirs(**selection) for source in selected: @@ -103,14 +119,17 @@ class StagedPlugin: raise ValueError("The updated plugin changed its installed name; reinstall it explicitly.") self.staged_digest = tree_digest(self.staged) self.metadata = self.target.parent / ".install-metadata.json" - self.previous = _bytes(self.metadata) - current = json.loads(self.previous) if self.previous is not None else {} - if current != plugin["old_metadata"]: + previous = _bytes(self.metadata) + current = _metadata_records(previous) + self.old_record = plugin["old_metadata"].get(self.target.name) + if current.get(self.target.name) != self.old_record: raise ValueError("Plugin install metadata changed while preparing the update; retry.") + self.new_record = plugin["new_metadata"].get(self.target.name) + if self.new_record is None: + raise ValueError("Plugin publication omitted its install metadata record.") self.target_digest = tree_digest(self.target) if self.target.exists() else None if self.target_digest != plugin["target_digest"]: raise ValueError("Plugin files changed while preparing the update; retry.") - self.proposed = (json.dumps(plugin["new_metadata"], indent=2, sort_keys=True) + "\n").encode() sources = member_sources(enabled_plugin_dirs(installing=self.target)) self.active = self.target.resolve() in sources self.members = {} @@ -123,26 +142,34 @@ class StagedPlugin: def publish(self, project: Path) -> None: import os import uuid + from hermes_cli.auth import _file_lock from pm.store import tree_digest if selection_snapshot() != self.configs: raise ValueError("Plugin enablement changed while preparing the update; retry.") if tree_digest(self.staged) != self.staged_digest: raise ValueError("Staged plugin files changed while preparing the update; retry.") - if _bytes(self.metadata) != self.previous: - raise ValueError("Plugin install metadata changed while preparing the update; retry.") - current = tree_digest(self.target) if self.target.exists() else None - if current != self.target_digest: - raise ValueError("Plugin files changed while preparing the update; retry.") - backup = self.target.parent / f".previous-{uuid.uuid4().hex}" - row = { - "kind": "plugin", "target": str(self.target), "backup": str(backup), "metadata": str(self.metadata), - "target_existed": self.target.exists(), "facts_before": _digest(runtime_facts_path(project)), - "metadata_before": base64.b64encode(self.previous).decode() if self.previous is not None else None, - "metadata_after": base64.b64encode(self.proposed).decode(), - } - _atomic_bytes(install_state_dir(project) / "publication.json", json.dumps(row).encode()) - if self.target.exists(): - os.replace(self.target, backup) - os.replace(self.staged, self.target) - _atomic_bytes(self.metadata, self.proposed) + lock = self.metadata.with_name(f"{self.metadata.name}.lock") + with _file_lock(lock, _METADATA_LOCK_HOLDER, 10.0, + "Timed out waiting for the plugin install metadata lock"): + previous = _bytes(self.metadata) + metadata = _metadata_records(previous) + if metadata.get(self.target.name) != self.old_record: + raise ValueError("Plugin install metadata changed while preparing the update; retry.") + metadata[self.target.name] = self.new_record + proposed = (json.dumps(metadata, indent=2, sort_keys=True) + "\n").encode() + current = tree_digest(self.target) if self.target.exists() else None + if current != self.target_digest: + raise ValueError("Plugin files changed while preparing the update; retry.") + backup = self.target.parent / f".previous-{uuid.uuid4().hex}" + row = { + "kind": "plugin", "target": str(self.target), "backup": str(backup), "metadata": str(self.metadata), + "target_existed": self.target.exists(), "facts_before": _digest(runtime_facts_path(project)), + "metadata_before": base64.b64encode(previous).decode() if previous is not None else None, + "metadata_after": base64.b64encode(proposed).decode(), + } + _atomic_bytes(install_state_dir(project) / "publication.json", json.dumps(row).encode()) + if self.target.exists(): + os.replace(self.target, backup) + os.replace(self.staged, self.target) + _atomic_bytes(self.metadata, proposed) diff --git a/tests/pm/test_worker_publication.py b/tests/pm/test_worker_publication.py index 25e49a51b9..6cfb3ebb03 100644 --- a/tests/pm/test_worker_publication.py +++ b/tests/pm/test_worker_publication.py @@ -304,6 +304,49 @@ def test_staged_publication_refuses_concurrent_input_edits(client, tmp_path, mon assert not (install_state_dir(repo) / "publication.json").exists() +def test_staged_publication_preserves_a_concurrent_sibling_install_record( + client, tmp_path, monkeypatch, isolated_python, +): + from pm.store import tree_digest + + _current_environment(tmp_path, monkeypatch, []) + home = tmp_path / "home" + target = home / "plugins/example" + target.mkdir(parents=True) + (target / "plugin.yaml").write_text("name: example\n") + (target / "code.py").write_text("old code") + (home / "config.yaml").write_text("plugins:\n enabled: [example]\n") + metadata = target.parent / ".install-metadata.json" + metadata.write_text('{"example":{"revision":"old"}}\n') + staged = tmp_path / "staged" + staged.mkdir() + (staged / "plugin.yaml").write_text("name: example\npython_dependencies: [fixture-dep==1]\n") + (staged / "code.py").write_text("new code") + concurrent = {"example": {"revision": "old"}, "sibling": {"revision": "sibling-new"}} + worker_toolchain( + client, + monkeypatch, + isolated_python, + "import json\nfrom pm.packages import Venv\n" + "def apply(*args, **kwargs):\n" + f" Path({str(metadata)!r}).write_text(json.dumps({concurrent!r}) + '\\n')\n" + " return {}\n" + "Venv.apply = apply\n", + ) + + client.sync_venv(explicit=True, staged_plugin={ + "target": str(target), "staged": str(staged), "target_digest": tree_digest(target), + "old_metadata": {"example": {"revision": "old"}}, + "new_metadata": {"example": {"revision": "new"}}, + }) + + assert json.loads(metadata.read_text()) == { + "example": {"revision": "new"}, + "sibling": {"revision": "sibling-new"}, + } + assert (target / "code.py").read_text() == "new code" + + def test_inactive_portable_publication_does_not_inspect_unrelated_dependency_manifests(client, tmp_path, monkeypatch): from pm.store import tree_digest _current_environment(tmp_path, monkeypatch, []) diff --git a/tools/mcp_tool_registration.py b/tools/mcp_tool_registration.py index 2170d8120a..5a16e82201 100644 --- a/tools/mcp_tool_registration.py +++ b/tools/mcp_tool_registration.py @@ -10,6 +10,7 @@ from dataclasses import dataclass from types import SimpleNamespace from typing import TYPE_CHECKING, Any, Callable, Dict, Iterable, List, Optional from tools.mcp_tool_common import _parse_boolish, _core, _resolve_tool_timeout, mcp_field, mcp_server_enabled +from tools import mcp_tool_config as _config from tools import mcp_tool_handlers as _handlers from tools import mcp_tool_schema as _schema from tools.mcp_tool_handlers import ( @@ -461,13 +462,24 @@ def _register_connected_into_current_scope(servers: dict) -> int: if scope is None: return 0 + # Callers that connect a subset (plugin go-live, a connector, orphan re-registration) pass only + # those names. A name they omit is judged against this profile's own config, or connecting one + # server would strip every other server's tools from the profile while their connections live on. + with _core._lock: + omitted = {_key_name(key) for key, scopes in _core._server_tool_scopes.items() + if scope in scopes and _key_name(key) not in servers} + profile_servers = _config._load_mcp_config() if omitted else {} + with _core._lock: stale = [] for key, scopes in _core._server_tool_scopes.items(): if scope not in scopes: continue + name = _key_name(key) + if name not in servers and name not in omitted: + continue # attached after the config read; the next pass judges it server = _core._servers.get(key) - config = servers.get(_key_name(key)) + config = servers[name] if name in servers else profile_servers.get(name) cross_profile = _key_scope(key) != scope if (config is None or not mcp_server_enabled(config) or server is None or getattr(server, "session", None) is None diff --git a/tools/tool_search.py b/tools/tool_search.py index 5c5a6ec7fd..6466dc49ad 100644 --- a/tools/tool_search.py +++ b/tools/tool_search.py @@ -238,7 +238,11 @@ def _search_description(deferred_count: int, listing: Optional[str], listing_for (f"Search {deferred_count} additional tools that are loaded on demand. " if deferred_count else "Search remote connector tools (email, calendars, issue trackers, and more). ") + "Takes a list of queries searched in parallel against the same " - "catalog; send one query per distinct capability you need. Returns " + "catalog; send one query per distinct capability you need. Queries are " + "keyword searches, not questions: the app or service name plus an action " + "and object, no filler words (`gmail send email`, `nvidia driver status`, " + "not `what's my GPU driver version?`); a word no tool contains makes the " + "query return nothing. Returns " "matching tool names grouped per query plus a shared map with each " "tool's description. Follow with " f"`{TOOL_DESCRIBE_NAME}` to load full parameter schemas, " @@ -280,7 +284,7 @@ def bridge_tool_schemas(deferred_count: int, listing: Optional[str] = None, "queries": { "type": "array", "items": {"type": "string"}, - "description": "Search queries, each a few keywords describing one capability (e.g. ['create github issue', 'send slack message']). Searched in parallel; results come back grouped per query. A single string is accepted and treated as one query.", + "description": "Keyword queries, one per capability: app or service name + action + object (e.g. ['github create issue', 'slack send message', 'gmail fetch emails']). Not questions or sentences: every word must appear in tool text, or the query returns nothing. Searched in parallel; results come back grouped per query. A single string is accepted and treated as one query.", }, "limit": { "type": "integer",