Merge remote-tracking branch 'origin/main' into ethie/pm-clean
# Conflicts: # hermes_cli/plugins_cmd.py # hermes_cli/plugins_cmd_catalog.py
This commit is contained in:
@@ -264,6 +264,10 @@ export function ConnectorOffer({ owner, request }: ConnectorOfferProps) {
|
||||
const { t } = useI18n()
|
||||
const copy = t.connectors
|
||||
const [reissuing, setReissuing] = useState<ReadonlySet<string>>(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<ReadonlySet<string>>(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<HTMLDivElement | null>(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<void> => {
|
||||
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)}
|
||||
/>
|
||||
)
|
||||
})}
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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<null | Readonly<Record<string, string>>> {
|
||||
const history = await requestGatewayForAgent<Partial<SessionHistoryResult>>(
|
||||
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
|
||||
}
|
||||
@@ -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
|
||||
|
||||
@@ -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]
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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.
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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, [])
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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",
|
||||
|
||||
Reference in New Issue
Block a user