feat(webui): add skills marketplace

This commit is contained in:
Xubin Ren
2026-07-30 01:13:37 +08:00
parent 129b74b4cf
commit ba0ba4749d
23 changed files with 3042 additions and 246 deletions
+18 -6
View File
@@ -70,6 +70,7 @@ class ContextBuilder:
def build_system_prompt(
self,
*,
active_skill_names: Sequence[str] | None = None,
channel: str | None = None,
session_summary: str | None = None,
workspace: Path | None = None,
@@ -91,13 +92,18 @@ class ContextBuilder:
if memory and not self._is_template_content(memory, "memory/MEMORY.md"):
parts.append(f"# Memory\n\n## Long-term Memory\n{memory}")
always_skills = self.skills.get_always_skills()
if always_skills:
always_content = self.skills.load_skills_for_context(always_skills)
if always_content:
parts.append(f"# Active Skills\n\n{always_content}")
active_skills = self.skills.get_always_skills()
active_skills.extend(
name
for name in (active_skill_names or ())
if name not in active_skills
)
if active_skills:
active_content = self.skills.load_skills_for_context(active_skills)
if active_content:
parts.append(f"# Active Skills\n\n{active_content}")
skills_summary = self.skills.build_skills_summary(exclude=set(always_skills))
skills_summary = self.skills.build_skills_summary(exclude=set(active_skills))
if skills_summary:
parts.append(render_template("agent/skills_section.md", skills_summary=skills_summary))
@@ -214,6 +220,11 @@ class ContextBuilder:
) -> list[dict[str, Any]]:
"""Build the complete message list for an LLM call."""
root = workspace or self.workspace
active_skill_names = (
self.skills.get_explicitly_invoked_skills(current_message)
if current_role == "user"
else []
)
user_content = self.build_user_content(current_message, image_paths=media)
blocks = list(runtime_context_blocks or ()) if current_role == "user" else []
merged, runtime_context_meta = append_runtime_context(user_content, blocks)
@@ -221,6 +232,7 @@ class ContextBuilder:
{
"role": "system",
"content": self.build_system_prompt(
active_skill_names=active_skill_names,
channel=channel,
session_summary=session_summary,
workspace=root,
+16
View File
@@ -17,6 +17,7 @@ _STRIP_SKILL_FRONTMATTER = re.compile(
r"^---\s*\r?\n(.*?)\r?\n---\s*\r?\n?",
re.DOTALL,
)
_SKILL_REFERENCE = re.compile(r"(?<![\w$])\$([A-Za-z0-9_-]+)")
class SkillsLoader:
@@ -109,6 +110,21 @@ class SkillsLoader:
]
return "\n\n---\n\n".join(parts)
def get_explicitly_invoked_skills(self, text: str) -> list[str]:
"""Resolve ``$skill-name`` references to enabled, available skills."""
if not text:
return []
available = {
entry["name"]
for entry in self.list_skills(filter_unavailable=True)
}
invoked: list[str] = []
for match in _SKILL_REFERENCE.finditer(text):
name = match.group(1)
if name in available and name not in invoked:
invoked.append(name)
return invoked
def build_skills_summary(self, exclude: set[str] | None = None) -> str:
"""
Build a summary of all skills (name, description, path, availability).
+3
View File
@@ -100,6 +100,7 @@ class ChannelManager:
webui_static_dist: bool = True,
webui_runtime_surface: str = "browser",
webui_runtime_capabilities: dict[str, Any] | None = None,
webui_skill_state_action: Callable[[set[str]], None] | None = None,
):
self.config = config
self.bus = bus
@@ -112,6 +113,7 @@ class ChannelManager:
self._webui_static_dist = webui_static_dist
self._webui_runtime_surface = webui_runtime_surface
self._webui_runtime_capabilities = dict(webui_runtime_capabilities or {})
self._webui_skill_state_action = webui_skill_state_action
self.channels: dict[str, BaseChannel] = {}
self._channel_owners: dict[str, str] = {}
self._channel_runtime_specs: dict[str, tuple[str, str]] = {}
@@ -178,6 +180,7 @@ class ChannelManager:
local_trigger_pending_ids=self._webui_local_trigger_pending_ids,
channel_feature_action=self.apply_channel_feature_action,
channel_runtime_status=self.get_status,
skill_state_action=self._webui_skill_state_action,
logger=logger,
)
kwargs["gateway"] = gateway
@@ -519,6 +519,8 @@ async def test_webui_skills_route_requires_token_and_hides_paths(
"name": "workspace-skill",
"description": "Workspace skill.",
"source": "workspace",
"enabled": True,
"deletable": True,
"available": True,
"unavailable_reason": "",
}
@@ -548,6 +550,212 @@ async def test_webui_skills_route_requires_token_and_hides_paths(
await server_task
@pytest.mark.asyncio
async def test_webui_skill_management_routes(
bus: MagicMock,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
skill_dir = tmp_path / "skills" / "custom-skill"
skill_dir.mkdir(parents=True)
(skill_dir / "SKILL.md").write_text(
"---\nname: custom-skill\ndescription: Custom skill.\n---\n",
encoding="utf-8",
)
def set_enabled(
workspace: Path,
name: str,
*,
enabled: bool,
disabled_skills: set[str],
) -> dict[str, Any]:
assert workspace == tmp_path
assert name == "custom-skill"
assert enabled is False
disabled_skills.add(name)
return {"name": name, "enabled": enabled, "deleted": False}
def delete(
workspace: Path,
name: str,
*,
disabled_skills: set[str],
) -> dict[str, Any]:
assert workspace == tmp_path
assert name == "custom-skill"
disabled_skills.discard(name)
for child in skill_dir.iterdir():
child.unlink()
skill_dir.rmdir()
return {"name": name, "enabled": False, "deleted": True}
monkeypatch.setattr("nanobot.webui.ws_http.set_webui_skill_enabled", set_enabled)
monkeypatch.setattr("nanobot.webui.ws_http.delete_webui_skill", delete)
port = _free_port()
channel = _ch(
bus,
session_manager=_seed_session(tmp_path),
workspace_path=tmp_path,
port=port,
)
server_task = asyncio.create_task(channel.start())
try:
token = channel.gateway.tokens.issue_api_token(300)
headers = {"Authorization": f"Bearer {token}"}
update_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/update"
"?name=custom-skill&enabled=false",
headers=headers,
)
assert update_response.status_code == 200
assert update_response.json()["last_action"]["enabled"] is False
custom = next(
item
for item in update_response.json()["skills"]
if item["name"] == "custom-skill"
)
assert custom["enabled"] is False
delete_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/delete"
"?name=custom-skill",
headers=headers,
)
assert delete_response.status_code == 200
assert delete_response.json()["last_action"]["deleted"] is True
assert all(
item["name"] != "custom-skill"
for item in delete_response.json()["skills"]
)
finally:
await channel.stop()
await server_task
@pytest.mark.asyncio
async def test_webui_skills_marketplace_routes_search_and_install(
bus: MagicMock,
tmp_path: Path,
monkeypatch: pytest.MonkeyPatch,
) -> None:
search = AsyncMock(return_value={
"query": "react",
"install_supported": True,
"skills": [{
"id": "acme/agent-skills/react-testing",
"skill_id": "react-testing",
"name": "React Testing",
"source": "acme/agent-skills",
"installs": 42,
"url": "https://skills.sh/acme/agent-skills/react-testing",
"installed": False,
}],
})
trending = AsyncMock(return_value={
"period": "24h",
"install_supported": True,
"skills": [{
"id": "acme/agent-skills/react-testing",
"skill_id": "react-testing",
"name": "React Testing",
"source": "acme/agent-skills",
"installs": 12,
"url": "https://skills.sh/acme/agent-skills/react-testing",
"installed": False,
"rank": 1,
}],
})
trends = AsyncMock(return_value={
"trends": {"acme/agent-skills/react-testing": [2, 4, 3, 8]},
})
async def install(source: str, skill_id: str, workspace: Path) -> dict[str, Any]:
assert source == "acme/agent-skills"
assert skill_id == "react-testing"
assert workspace == tmp_path
skill_dir = workspace / "skills" / skill_id
skill_dir.mkdir(parents=True)
(skill_dir / "SKILL.md").write_text(
"---\nname: react-testing\ndescription: Test React apps.\n---\n",
encoding="utf-8",
)
return {"installed": True, "already_installed": False, "name": skill_id}
install_mock = AsyncMock(side_effect=install)
monkeypatch.setattr("nanobot.webui.ws_http.search_marketplace_skills", search)
monkeypatch.setattr("nanobot.webui.ws_http.trending_marketplace_skills", trending)
monkeypatch.setattr("nanobot.webui.ws_http.marketplace_skill_trends", trends)
monkeypatch.setattr("nanobot.webui.ws_http.install_marketplace_skill", install_mock)
port = _free_port()
channel = _ch(
bus,
session_manager=_seed_session(tmp_path),
workspace_path=tmp_path,
port=port,
)
server_task = asyncio.create_task(channel.start())
try:
denied = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/search?q=react"
)
assert denied.status_code == 401
token = channel.gateway.tokens.issue_api_token(300)
headers = {"Authorization": f"Bearer {token}"}
search_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/search?q=react",
headers=headers,
)
assert search_response.status_code == 200
assert search_response.json()["skills"][0]["skill_id"] == "react-testing"
search.assert_awaited_once_with("react", tmp_path)
trending_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/trending",
headers=headers,
)
assert trending_response.status_code == 200
assert trending_response.json()["period"] == "24h"
trending.assert_awaited_once_with(tmp_path)
trends_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/trends"
"?id=acme%2Fagent-skills%2Freact-testing",
headers=headers,
)
assert trends_response.status_code == 200
assert trends_response.json()["trends"] == {
"acme/agent-skills/react-testing": [2, 4, 3, 8],
}
trends.assert_awaited_once_with(["acme/agent-skills/react-testing"])
params = urlencode({
"source": "acme/agent-skills",
"skill": "react-testing",
})
install_response = await _http_get(
f"http://127.0.0.1:{port}/api/webui/skills/install?{params}",
headers=headers,
)
assert install_response.status_code == 200
body = install_response.json()
assert body["last_action"] == {
"installed": True,
"already_installed": False,
"name": "react-testing",
}
assert next(
skill for skill in body["skills"] if skill["name"] == "react-testing"
)["source"] == "workspace"
install_mock.assert_awaited_once()
finally:
await channel.stop()
await server_task
@pytest.mark.asyncio
async def test_cli_apps_routes_require_token_and_return_payload(
bus: MagicMock,
+6
View File
@@ -2139,6 +2139,11 @@ def _run_gateway(
def _webui_runtime_model_name() -> str | None:
return agent.model.strip() or None
def _webui_skill_state_action(disabled_skills: set[str]) -> None:
config.agents.defaults.disabled_skills = sorted(disabled_skills)
agent.context.skills.disabled_skills = set(disabled_skills)
agent.subagents.disabled_skills = set(disabled_skills)
# Create channel manager (forwards SessionManager so the WebSocket channel
# can serve the embedded webui's REST surface).
channels = ChannelManager(
@@ -2153,6 +2158,7 @@ def _run_gateway(
webui_static_dist=webui_static_dist,
webui_runtime_surface=webui_runtime_surface,
webui_runtime_capabilities=webui_runtime_capabilities,
webui_skill_state_action=_webui_skill_state_action,
)
def _pick_heartbeat_target() -> tuple[str, str]:
+2
View File
@@ -58,6 +58,7 @@ def build_gateway_services(
local_trigger_pending_ids: Callable[[str], set[str]] | None = None,
channel_feature_action: Callable[..., Any] | None = None,
channel_runtime_status: Callable[[], dict[str, Any]] | None = None,
skill_state_action: Callable[[set[str]], None] | None = None,
logger: Any = default_logger,
) -> GatewayServices:
tokens = GatewayTokenStore()
@@ -101,6 +102,7 @@ def build_gateway_services(
local_trigger_pending_ids=local_trigger_pending_ids,
channel_feature_action=channel_feature_action,
channel_runtime_status=channel_runtime_status,
skill_state_action=skill_state_action,
log=logger,
)
return GatewayServices(
+161 -7
View File
@@ -2,10 +2,23 @@
from __future__ import annotations
import json
import shlex
import shutil
from pathlib import Path
from typing import Any
from nanobot.agent.skills import SkillsLoader
from nanobot.config.loader import load_config, save_config
class SkillManagementError(Exception):
"""A safe skill-management error for the WebUI."""
def __init__(self, message: str, *, status: int = 400) -> None:
super().__init__(message)
self.message = message
self.status = status
def webui_skills_payload(
@@ -14,12 +27,17 @@ def webui_skills_payload(
disabled_skills: set[str] | None = None,
) -> dict[str, Any]:
"""Return agent skills without leaking local filesystem paths."""
loader = SkillsLoader(workspace_path, disabled_skills=disabled_skills)
loader = SkillsLoader(workspace_path)
entries = sorted(
loader.list_skills(filter_unavailable=False),
key=lambda entry: (entry.get("source") != "workspace", entry["name"]),
)
return {"skills": [_skill_payload(loader, entry) for entry in entries]}
return {
"skills": [
_skill_payload(loader, entry, disabled_skills=disabled_skills)
for entry in entries
]
}
def webui_skill_detail_payload(
@@ -29,26 +47,114 @@ def webui_skill_detail_payload(
disabled_skills: set[str] | None = None,
) -> dict[str, Any] | None:
"""Return a single skill's safe detail payload."""
loader = SkillsLoader(workspace_path, disabled_skills=disabled_skills)
loader = SkillsLoader(workspace_path)
entries = loader.list_skills(filter_unavailable=False)
entry = next((item for item in entries if item["name"] == name), None)
if entry is None:
return None
metadata = loader.get_skill_metadata(name)
return {
**_skill_payload(loader, entry),
**_skill_payload(
loader,
entry,
metadata=metadata,
disabled_skills=disabled_skills,
),
"requirements": loader.get_skill_requirements(name),
"install_options": _install_options(metadata),
"raw_markdown": loader.load_skill(name) or "",
}
def _skill_payload(loader: SkillsLoader, entry: dict[str, str]) -> dict[str, Any]:
def set_webui_skill_enabled(
workspace_path: Path,
name: str,
*,
enabled: bool,
disabled_skills: set[str],
) -> dict[str, Any]:
"""Persist and apply one skill's enabled state."""
_require_skill_entry(workspace_path, name)
config = load_config()
next_disabled = set(config.agents.defaults.disabled_skills)
if enabled:
next_disabled.discard(name)
else:
next_disabled.add(name)
if next_disabled != set(config.agents.defaults.disabled_skills):
config.agents.defaults.disabled_skills = sorted(next_disabled)
save_config(config)
disabled_skills.clear()
disabled_skills.update(next_disabled)
return {"name": name, "enabled": enabled, "deleted": False}
def delete_webui_skill(
workspace_path: Path,
name: str,
*,
disabled_skills: set[str],
) -> dict[str, Any]:
"""Delete one workspace skill and remove its disabled-state entry."""
entry = _require_skill_entry(workspace_path, name)
if entry.get("source") != "workspace":
raise SkillManagementError("built-in skills cannot be deleted", status=403)
skills_root = (workspace_path.expanduser().resolve() / "skills").resolve()
target = skills_root / name
if target.parent != skills_root:
raise SkillManagementError("invalid skill name")
if target.is_symlink():
target.unlink()
elif target.is_dir():
shutil.rmtree(target)
else:
raise SkillManagementError("skill directory was not found", status=404)
config = load_config()
next_disabled = set(config.agents.defaults.disabled_skills)
if name in next_disabled:
next_disabled.remove(name)
config.agents.defaults.disabled_skills = sorted(next_disabled)
save_config(config)
disabled_skills.clear()
disabled_skills.update(next_disabled)
return {"name": name, "enabled": False, "deleted": True}
def _require_skill_entry(workspace_path: Path, name: str) -> dict[str, str]:
if not name or "/" in name or "\\" in name:
raise SkillManagementError("invalid skill name")
entry = next(
(
item
for item in SkillsLoader(workspace_path).list_skills(filter_unavailable=False)
if item["name"] == name
),
None,
)
if entry is None:
raise SkillManagementError("skill not found", status=404)
return entry
def _skill_payload(
loader: SkillsLoader,
entry: dict[str, str],
*,
metadata: dict[str, Any] | None = None,
disabled_skills: set[str] | None = None,
) -> dict[str, Any]:
name = entry["name"]
metadata = loader.get_skill_metadata(name)
metadata = metadata if metadata is not None else loader.get_skill_metadata(name)
available, unavailable_reason = loader.get_skill_availability(name)
source = entry.get("source", "unknown")
return {
"name": name,
"description": _description(metadata, name),
"source": entry.get("source", "unknown"),
"source": source,
"enabled": name not in (disabled_skills or set()),
"deletable": source == "workspace",
"available": available,
"unavailable_reason": unavailable_reason,
}
@@ -59,3 +165,51 @@ def _description(metadata: dict[str, Any] | None, fallback: str) -> str:
return fallback
value = metadata.get("description")
return value.strip() if isinstance(value, str) and value.strip() else fallback
def _nanobot_metadata(metadata: dict[str, Any] | None) -> dict[str, Any]:
if metadata is None:
return {}
raw = metadata.get("metadata")
if isinstance(raw, str):
try:
raw = json.loads(raw)
except (json.JSONDecodeError, TypeError):
return {}
if not isinstance(raw, dict):
return {}
payload = raw.get("nanobot", raw.get("openclaw", {}))
return payload if isinstance(payload, dict) else {}
def _install_options(metadata: dict[str, Any] | None) -> list[dict[str, str]]:
"""Return safe, copyable setup commands declared by a skill."""
install = _nanobot_metadata(metadata).get("install")
if not isinstance(install, list):
return []
options: list[dict[str, str]] = []
for item in install:
if not isinstance(item, dict):
continue
kind = item.get("kind")
package = item.get("formula") if kind == "brew" else item.get("package")
if not isinstance(package, str) or not package.strip():
continue
if kind == "brew":
command = f"brew install {shlex.quote(package.strip())}"
elif kind == "apt":
command = f"sudo apt-get install -y {shlex.quote(package.strip())}"
else:
continue
option_id = item.get("id")
label = item.get("label")
options.append(
{
"id": option_id if isinstance(option_id, str) else kind,
"kind": kind,
"label": label if isinstance(label, str) else f"Install with {kind}",
"command": command,
}
)
return options
+395
View File
@@ -0,0 +1,395 @@
"""Search and install skills from the skills.sh catalog."""
from __future__ import annotations
import asyncio
import os
import re
import shutil
import time
from pathlib import Path
from typing import Any
import httpx
from nanobot.agent.skills import SkillsLoader
from nanobot.security.network import PinnedDNSAsyncTransport
_SEARCH_URL = "https://skills.sh/api/search"
_TRENDING_URL = "https://skills.sh/api/skills/trending/0"
_SKILL_PAGE_BASE_URL = "https://www.skills.sh"
_ALL_TIME_URLS = (
"https://skills.sh/api/skills/all-time/0",
"https://skills.sh/api/skills/all-time/1",
)
_TREND_VALUES_RE = re.compile(r'\\"values\\":\s*\[([0-9,\s]+)\]')
_SOURCE_RE = re.compile(
r"^[A-Za-z0-9](?:[A-Za-z0-9-]{0,37}[A-Za-z0-9])?/"
r"[A-Za-z0-9](?:[A-Za-z0-9_.-]{0,98}[A-Za-z0-9])?$"
)
_SKILL_RE = re.compile(r"^[a-z0-9]+(?:-[a-z0-9]+)*$")
_ANSI_RE = re.compile(r"\x1b\[[0-?]*[ -/]*[@-~]")
_INSTALL_TIMEOUT_SECONDS = 120
_WEEKLY_CACHE_TTL_SECONDS = 300
# The skills CLI's OpenClaw adapter copies into <workspace>/skills, nanobot's layout too.
_CLI_AGENT = "openclaw"
_weekly_cache: dict[tuple[str, str], list[int]] = {}
_weekly_cache_expires_at = 0.0
class SkillsMarketplaceError(Exception):
"""A safe error that can be returned to the WebUI."""
def __init__(self, message: str, *, status: int = 400) -> None:
super().__init__(message)
self.message = message
self.status = status
def skills_install_supported() -> bool:
"""Return whether the official skills CLI can be launched."""
return shutil.which("npx") is not None
async def trending_marketplace_skills(
workspace_path: Path,
*,
limit: int = 8,
) -> dict[str, Any]:
"""Return a source-diverse snapshot of skills.sh's real 24-hour leaderboard."""
try:
async with _skills_client() as client:
response = await client.get(_TRENDING_URL)
response.raise_for_status()
payload = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise SkillsMarketplaceError(
"skills.sh trending skills are temporarily unavailable",
status=502,
) from exc
installed = _installed_skill_names(workspace_path)
rows = payload.get("skills", []) if isinstance(payload, dict) else []
skills: list[dict[str, Any]] = []
seen_sources: set[str] = set()
for rank, row in enumerate(rows, start=1):
if not isinstance(row, dict):
continue
source = row.get("source")
if not isinstance(source, str) or source in seen_sources:
continue
skill = _marketplace_skill(row, installed, rank=rank)
if skill is None:
continue
seen_sources.add(source)
skills.append(skill)
if len(skills) >= min(max(limit, 1), 20):
break
return {
"skills": skills,
"period": "24h",
"install_supported": skills_install_supported(),
}
async def search_marketplace_skills(
query: str,
workspace_path: Path,
*,
limit: int = 20,
) -> dict[str, Any]:
"""Search skills.sh and annotate results already installed in this workspace."""
normalized = " ".join(query.split())
if len(normalized) < 2:
raise SkillsMarketplaceError("search query must contain at least 2 characters")
if len(normalized) > 100:
raise SkillsMarketplaceError("search query is too long")
try:
async with _skills_client() as client:
response = await client.get(
_SEARCH_URL,
params={"q": normalized, "limit": min(max(limit, 1), 50)},
)
response.raise_for_status()
payload = response.json()
except (httpx.HTTPError, ValueError) as exc:
raise SkillsMarketplaceError(
"skills.sh search is temporarily unavailable",
status=502,
) from exc
installed = _installed_skill_names(workspace_path)
rows = payload.get("skills", []) if isinstance(payload, dict) else []
skills = []
for row in rows:
if not isinstance(row, dict):
continue
skill = _marketplace_skill(row, installed)
if skill is not None:
skills.append(skill)
return {
"query": normalized,
"skills": skills,
"install_supported": skills_install_supported(),
}
async def marketplace_skill_trends(
skill_ids: list[str] | None = None,
) -> dict[str, dict[str, list[int]]]:
"""Return install history independently, filling requested cache misses."""
requested = _valid_skill_refs(skill_ids or [])
async with _skills_client() as client:
weekly_installs = await _load_weekly_installs(client)
missing = [ref for ref in requested if ref not in weekly_installs]
if missing:
weekly_installs.update(await _load_skill_page_trends(client, missing))
selected = requested or list(weekly_installs)
return {
"trends": {
f"{source}/{skill_id}": values
for source, skill_id in selected
if (values := weekly_installs.get((source, skill_id))) is not None
}
}
async def install_marketplace_skill(
source: str,
skill_id: str,
workspace_path: Path,
) -> dict[str, Any]:
"""Install one skills.sh result into ``<workspace>/skills``."""
if not _SOURCE_RE.fullmatch(source):
raise SkillsMarketplaceError("invalid skill source")
if not _valid_skill_id(skill_id):
raise SkillsMarketplaceError("invalid skill name")
loader = SkillsLoader(workspace_path)
existing = {
entry["name"]: entry
for entry in loader.list_skills(filter_unavailable=False)
}
if skill_id in existing:
return {"installed": True, "already_installed": True, "name": skill_id}
npx = shutil.which("npx")
if npx is None:
raise SkillsMarketplaceError(
"Node.js with npx is required to install skills",
status=503,
)
workspace = workspace_path.expanduser().resolve()
workspace.mkdir(parents=True, exist_ok=True)
env = os.environ.copy()
env["DISABLE_TELEMETRY"] = "1"
command = (
npx,
"--yes",
"skills@latest",
"add",
source,
"--skill",
skill_id,
"--agent",
_CLI_AGENT,
"--copy",
"--yes",
)
process = await asyncio.create_subprocess_exec(
*command,
cwd=str(workspace),
env=env,
stdout=asyncio.subprocess.PIPE,
stderr=asyncio.subprocess.STDOUT,
)
try:
output, _ = await asyncio.wait_for(
process.communicate(),
timeout=_INSTALL_TIMEOUT_SECONDS,
)
except TimeoutError as exc:
process.kill()
await process.communicate()
raise SkillsMarketplaceError("skill installation timed out", status=504) from exc
if process.returncode != 0:
detail = _safe_output_tail(output)
message = "skill installation failed"
if detail:
message = f"{message}: {detail}"
raise SkillsMarketplaceError(message, status=502)
installed = next(
(
entry
for entry in loader.list_skills(filter_unavailable=False)
if entry["source"] == "workspace" and entry["name"] == skill_id
),
None,
)
if installed is None:
raise SkillsMarketplaceError(
"installer completed but the skill was not found in this workspace",
status=502,
)
return {"installed": True, "already_installed": False, "name": skill_id}
def _skills_client() -> httpx.AsyncClient:
return httpx.AsyncClient(
transport=PinnedDNSAsyncTransport(),
timeout=10.0,
follow_redirects=False,
)
def _installed_skill_names(workspace_path: Path) -> set[str]:
return {
entry["name"]
for entry in SkillsLoader(workspace_path).list_skills(filter_unavailable=False)
}
def _marketplace_skill(
row: dict[str, Any],
installed: set[str],
*,
rank: int | None = None,
) -> dict[str, Any] | None:
source = row.get("source")
skill_id = row.get("skillId")
if not isinstance(source, str) or not _SOURCE_RE.fullmatch(source):
return None
if not isinstance(skill_id, str) or not _valid_skill_id(skill_id):
return None
display_name = row.get("name")
if not isinstance(display_name, str) or not display_name.strip():
display_name = skill_id
installs = row.get("installs")
skill: dict[str, Any] = {
"id": f"{source}/{skill_id}",
"skill_id": skill_id,
"name": display_name.strip(),
"source": source,
"installs": installs if isinstance(installs, int) and installs >= 0 else 0,
"url": f"https://skills.sh/{source}/{skill_id}",
"installed": skill_id in installed,
}
if rank is not None:
skill["rank"] = rank
return skill
async def _load_weekly_installs(
client: httpx.AsyncClient,
) -> dict[tuple[str, str], list[int]]:
global _weekly_cache, _weekly_cache_expires_at
now = time.monotonic()
if now < _weekly_cache_expires_at:
return _weekly_cache
responses = await asyncio.gather(
*(client.get(url) for url in _ALL_TIME_URLS),
return_exceptions=True,
)
history: dict[tuple[str, str], list[int]] = {}
successful = False
for response in responses:
if isinstance(response, BaseException):
continue
try:
response.raise_for_status()
payload = response.json()
except (httpx.HTTPError, ValueError):
continue
successful = True
rows = payload.get("skills", []) if isinstance(payload, dict) else []
for row in rows:
if not isinstance(row, dict):
continue
source = row.get("source")
skill_id = row.get("skillId")
values = row.get("weeklyInstalls")
if (
isinstance(source, str)
and isinstance(skill_id, str)
and isinstance(values, list)
):
clean = [
value
for value in values
if isinstance(value, int) and not isinstance(value, bool) and value >= 0
]
if len(clean) >= 2:
history[(source, skill_id)] = clean
if successful:
_weekly_cache = history
_weekly_cache_expires_at = now + _WEEKLY_CACHE_TTL_SECONDS
return history
def _valid_skill_refs(skill_ids: list[str]) -> list[tuple[str, str]]:
refs: list[tuple[str, str]] = []
for value in skill_ids[:20]:
if not isinstance(value, str) or "/" not in value:
continue
source, skill_id = value.rsplit("/", 1)
ref = (source, skill_id)
if (
_SOURCE_RE.fullmatch(source)
and _valid_skill_id(skill_id)
and ref not in refs
):
refs.append(ref)
return refs
async def _load_skill_page_trends(
client: httpx.AsyncClient,
refs: list[tuple[str, str]],
) -> dict[tuple[str, str], list[int]]:
semaphore = asyncio.Semaphore(6)
async def fetch(ref: tuple[str, str]) -> tuple[tuple[str, str], list[int]]:
source, skill_id = ref
try:
async with semaphore:
response = await client.get(
f"{_SKILL_PAGE_BASE_URL}/{source}/{skill_id}"
)
response.raise_for_status()
except httpx.HTTPError:
return ref, []
match = _TREND_VALUES_RE.search(response.text)
if match is None:
return ref, []
values = [
int(value)
for value in match.group(1).split(",")
if value.strip()
]
return ref, values if len(values) >= 2 else []
return dict(await asyncio.gather(*(fetch(ref) for ref in refs)))
def _valid_skill_id(value: str) -> bool:
return len(value) <= 64 and _SKILL_RE.fullmatch(value) is not None
def _safe_output_tail(output: bytes | None) -> str:
if not output:
return ""
text = _ANSI_RE.sub("", output.decode("utf-8", errors="replace"))
lines = [line.strip() for line in text.splitlines() if line.strip()]
return " · ".join(lines[-3:])[-600:]
+165 -2
View File
@@ -87,7 +87,20 @@ from nanobot.webui.sidebar_state import (
read_webui_sidebar_state,
write_webui_sidebar_state,
)
from nanobot.webui.skills_api import webui_skill_detail_payload, webui_skills_payload
from nanobot.webui.skills_api import (
SkillManagementError,
delete_webui_skill,
set_webui_skill_enabled,
webui_skill_detail_payload,
webui_skills_payload,
)
from nanobot.webui.skills_marketplace import (
SkillsMarketplaceError,
install_marketplace_skill,
marketplace_skill_trends,
search_marketplace_skills,
trending_marketplace_skills,
)
from nanobot.webui.thread_disk import delete_webui_thread
from nanobot.webui.transcript import build_webui_thread_response
from nanobot.webui.workspaces import WebUIWorkspaceController
@@ -171,6 +184,7 @@ class GatewayHTTPHandler:
local_trigger_pending_ids: Callable[[str], set[str]] | None = None,
channel_feature_action: Callable[..., Any] | None = None,
channel_runtime_status: Callable[[], dict[str, Any]] | None = None,
skill_state_action: Callable[[set[str]], None] | None = None,
log: Any = logger,
) -> None:
self.config = config
@@ -183,7 +197,8 @@ class GatewayHTTPHandler:
self.ingress = ingress
self.workspaces = workspaces
self.skills_workspace_path = skills_workspace_path
self.disabled_skills = disabled_skills or set()
self.disabled_skills = disabled_skills if disabled_skills is not None else set()
self.skill_state_action = skill_state_action
self.cron_service = cron_service
self.local_trigger_store = local_trigger_store
self.cron_pending_job_ids = cron_pending_job_ids
@@ -795,6 +810,18 @@ class GatewayHTTPHandler:
return self._handle_commands(request)
if got == "/api/workspaces":
return self._handle_workspaces(connection, request)
if got == "/api/webui/skills/search":
return await self._handle_webui_skills_search(request)
if got == "/api/webui/skills/trending":
return await self._handle_webui_skills_trending(request)
if got == "/api/webui/skills/trends":
return await self._handle_webui_skill_trends(request)
if got == "/api/webui/skills/install":
return await self._handle_webui_skill_install(connection, request)
if got == "/api/webui/skills/update":
return self._handle_webui_skill_update(request)
if got == "/api/webui/skills/delete":
return self._handle_webui_skill_delete(connection, request)
if got == "/api/webui/skills":
return self._handle_webui_skills(request)
m = re.match(r"^/api/webui/skills/([^/]+)$", got)
@@ -830,6 +857,142 @@ class GatewayHTTPHandler:
)
)
async def _handle_webui_skills_search(self, request: WsRequest) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
query = _query_first(_parse_query(request.path), "q") or ""
try:
payload = await search_marketplace_skills(query, self.skills_workspace_path)
except SkillsMarketplaceError as exc:
return _http_error(exc.status, exc.message)
except Exception:
self._log.exception("skills.sh search failed")
return _http_error(500, "skills.sh search failed")
return _http_json_response(payload)
async def _handle_webui_skills_trending(self, request: WsRequest) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
try:
payload = await trending_marketplace_skills(self.skills_workspace_path)
except SkillsMarketplaceError as exc:
return _http_error(exc.status, exc.message)
except Exception:
self._log.exception("skills.sh trending lookup failed")
return _http_error(500, "skills.sh trending lookup failed")
return _http_json_response(payload)
async def _handle_webui_skill_trends(self, request: WsRequest) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
skill_ids = _parse_query(request.path).get("id", [])
try:
payload = await marketplace_skill_trends(skill_ids)
except Exception:
self._log.exception("skills.sh trend history lookup failed")
return _http_error(500, "skills.sh trend history lookup failed")
return _http_json_response(payload)
async def _handle_webui_skill_install(
self,
connection: Any,
request: WsRequest,
) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
if not self._allow_webui_package_install(connection, request):
return _http_error(403, "remote skill installation is disabled")
query = _parse_query(request.path)
source = _query_first(query, "source") or ""
skill_id = _query_first(query, "skill") or ""
try:
action = await install_marketplace_skill(
source,
skill_id,
self.skills_workspace_path,
)
except SkillsMarketplaceError as exc:
return _http_error(exc.status, exc.message)
except Exception:
self._log.exception("skill installation failed")
return _http_error(500, "skill installation failed")
return _http_json_response({
**webui_skills_payload(
self.skills_workspace_path,
disabled_skills=self.disabled_skills,
),
"last_action": action,
})
def _allow_webui_package_install(self, connection: Any, request: WsRequest) -> bool:
if _is_local_browser_request(connection, request.headers):
return True
try:
from nanobot.config.loader import load_config
return bool(load_config().tools.webui_allow_remote_package_install)
except Exception:
self._log.exception("failed to load remote package install policy")
return False
def _handle_webui_skill_update(self, request: WsRequest) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
query = _parse_query(request.path)
name = _query_first(query, "name") or ""
raw_enabled = (_query_first(query, "enabled") or "").lower()
if raw_enabled not in {"true", "false"}:
return _http_error(400, "enabled must be true or false")
try:
action = set_webui_skill_enabled(
self.skills_workspace_path,
name,
enabled=raw_enabled == "true",
disabled_skills=self.disabled_skills,
)
except SkillManagementError as exc:
return _http_error(exc.status, exc.message)
self._apply_skill_state()
return _http_json_response({
**webui_skills_payload(
self.skills_workspace_path,
disabled_skills=self.disabled_skills,
),
"last_action": action,
})
def _handle_webui_skill_delete(
self,
connection: Any,
request: WsRequest,
) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")
if not self._allow_webui_package_install(connection, request):
return _http_error(403, "remote skill deletion is disabled")
name = _query_first(_parse_query(request.path), "name") or ""
try:
action = delete_webui_skill(
self.skills_workspace_path,
name,
disabled_skills=self.disabled_skills,
)
except SkillManagementError as exc:
return _http_error(exc.status, exc.message)
self._apply_skill_state()
return _http_json_response({
**webui_skills_payload(
self.skills_workspace_path,
disabled_skills=self.disabled_skills,
),
"last_action": action,
})
def _apply_skill_state(self) -> None:
if self.skill_state_action is not None:
self.skill_state_action(set(self.disabled_skills))
def _handle_webui_skill_detail(self, request: WsRequest, raw_name: str) -> Response:
if not self.check_api_token(request):
return _http_error(401, "Unauthorized")