From 371b72ddd1afafa9265862de3b0a2fc1e7df0f42 Mon Sep 17 00:00:00 2001 From: ciregenz Date: Fri, 19 Jun 2026 03:32:16 -0700 Subject: [PATCH] [eric] settings-agent/skills: fix concurrent-write lost update (serialize), registry install slug clobber (dedup), dialog stale-response races, partial-write + skills.sh boundary hardening --- backend/apps/agents/settings_meta_server.py | 10 ++- backend/apps/skill_registry/skill_registry.py | 12 ++- backend/apps/skills/skills.py | 21 +++++ backend/main.py | 86 +++++++++++-------- .../tests/test_settings_meta_concurrency.py | 62 +++++++++++++ .../tests/test_skill_registry_community.py | 15 ++++ .../pages/Skills/CommunitySkillsDialog.tsx | 19 +++- 7 files changed, 180 insertions(+), 45 deletions(-) create mode 100644 backend/tests/test_settings_meta_concurrency.py diff --git a/backend/apps/agents/settings_meta_server.py b/backend/apps/agents/settings_meta_server.py index d630217e..63ee1816 100644 --- a/backend/apps/agents/settings_meta_server.py +++ b/backend/apps/agents/settings_meta_server.py @@ -108,12 +108,16 @@ def _format_read(settings: dict) -> str: def _format_write(outcomes: dict) -> str: applied = [f for f, o in outcomes.items() if o.get("status") == "applied"] - refused = {f: o for f, o in outcomes.items() if o.get("status") not in ("applied", None)} parts = [] if applied: parts.append("Applied: " + ", ".join(sorted(applied))) - for field, o in refused.items(): - parts.append(f"Refused {field}: {o.get('reason', o.get('status'))}") + for field, o in outcomes.items(): + status = o.get("status") + if status in ("applied", None): + continue + # "error" is transient (retryable); "refused"/"unknown" are not. + verb = "Failed" if status == "error" else "Refused" + parts.append(f"{verb} {field}: {o.get('reason', status)}") if not parts: return "No changes were applied." return "\n".join(parts) diff --git a/backend/apps/skill_registry/skill_registry.py b/backend/apps/skill_registry/skill_registry.py index ef99b10b..ab97c409 100644 --- a/backend/apps/skill_registry/skill_registry.py +++ b/backend/apps/skill_registry/skill_registry.py @@ -385,7 +385,10 @@ async def _community_search(q: str, limit: int) -> dict: skills = [] for s in (data.get("skills") or [])[:limit]: src = s.get("source", "") - installs = s.get("installs", 0) + try: + installs = int(s.get("installs") or 0) + except (TypeError, ValueError): + installs = 0 skills.append({ "name": s.get("name", ""), "description": f"{installs:,} installs", @@ -437,9 +440,12 @@ async def registry_install(req: _InstallRequest): if not req.confirm: return {"installed": False, "disclosure": disclosure} - from backend.apps.skills.skills import write_folder_skill + from backend.apps.skills.skills import write_folder_skill, unique_skill_slug + # Never clobber an existing local skill that happens to share this slug; a + # wild-registry name collision lands as a copy instead of overwriting. + slug = unique_skill_slug(resolved["skill_id"]) skill = write_folder_skill( - resolved["skill_id"], + slug, resolved["files"], {"name": resolved["name"], "description": resolved["description"]}, ) diff --git a/backend/apps/skills/skills.py b/backend/apps/skills/skills.py index 4ebef13e..49b8919e 100644 --- a/backend/apps/skills/skills.py +++ b/backend/apps/skills/skills.py @@ -278,6 +278,27 @@ def _safe_slug(raw: str) -> str: return slug or "skill" +def _skill_exists(slug: str) -> bool: + return ( + slug in _load_index() + or os.path.isfile(os.path.join(SKILLS_DIR, f"{slug}.md")) + or os.path.isdir(os.path.join(SKILLS_DIR, slug)) + ) + + +def unique_skill_slug(base: str) -> str: + """A free slug for `base`, suffixing -2, -3, ... on collision. Lets a + registry install land beside a same-named skill instead of silently + overwriting the user's existing one.""" + slug = _safe_slug(base) + if not _skill_exists(slug): + return slug + i = 2 + while _skill_exists(f"{slug}-{i}"): + i += 1 + return f"{slug}-{i}" + + def write_folder_skill(skill_id: str, files: dict[str, str], meta: dict) -> Skill: """Write a multi-file skill folder (relpath -> content) under SKILLS_DIR and index it. `files` must include a 'SKILL.md'. Shared by registry install and diff --git a/backend/main.py b/backend/main.py index 17201672..68aad3bc 100644 --- a/backend/main.py +++ b/backend/main.py @@ -742,6 +742,11 @@ async def mcp_meta(action: str, request: Request): return JSONResponse({"error": f"unknown action: {action}"}, status_code=400) +# Serializes agent-side SettingsWrite read-modify-writes so concurrent autonomous +# agents can't clobber each other's edits (see the lock's use below). +_settings_meta_write_lock = asyncio.Lock() + + @app.post("/api/settings-meta/{action}") async def settings_meta(action: str, request: Request): """Back the openswarm-settings-meta stdio MCP server (agent-editable Settings). @@ -779,46 +784,57 @@ async def settings_meta(action: str, request: Request): if not isinstance(changes, dict) or not changes: return JSONResponse({"error": "changes must be a non-empty object of field -> value"}, status_code=400) - settings = load_settings() - session = agent_manager.sessions.get(parent_session_id) if parent_session_id else None - if session is not None: - powering = resolve_powering_credential(session.model, settings) - else: - # No live session to anchor the guard: fail safe, protect every credential. - powering = PoweringCredential(kind="unknown", provider="unknown", label="this run") - valid_fields = set(AppSettings.model_fields.keys()) outcomes: dict[str, dict] = {} - staged: dict = {} - for field, value in changes.items(): - if field not in valid_fields: - outcomes[field] = {"status": "unknown", "reason": "not a settings field"} - elif field in SERVER_OWNED_FIELDS: - outcomes[field] = {"status": "refused", "reason": "managed by your subscription/connection; change it in the Subscription section"} - elif write_would_suicide(field, value, powering): - outcomes[field] = {"status": "refused", "reason": f"would disconnect {powering.label}, which is powering this run"} + # Serialize the read-modify-write: SettingsWrite goes through update_settings, + # which awaits (so two autonomous agents would interleave and clobber each + # other's fields while BOTH got an "applied" result). The lock makes agent + # writes serial so the last load always sees the prior write. (Agent vs the + # renderer's own PUT stays the pre-existing full-object-replace race.) + async with _settings_meta_write_lock: + settings = load_settings() + session = agent_manager.sessions.get(parent_session_id) if parent_session_id else None + if session is not None: + powering = resolve_powering_credential(session.model, settings) else: - staged[field] = value + # No live session to anchor the guard: fail safe, protect every credential. + powering = PoweringCredential(kind="unknown", provider="unknown", label="this run") - if staged: - merged = settings.model_dump() - merged.update(staged) - try: - new_body = AppSettings(**merged) - except ValidationError as e: - bad = {str(err["loc"][0]) for err in e.errors() if err.get("loc")} - for f in bad & set(staged.keys()): - outcomes[f] = {"status": "refused", "reason": "invalid value for this field"} - staged.pop(f, None) - new_body = None - if staged: - merged = settings.model_dump() - merged.update(staged) + staged: dict = {} + for field, value in changes.items(): + if field not in valid_fields: + outcomes[field] = {"status": "unknown", "reason": "not a settings field"} + elif field in SERVER_OWNED_FIELDS: + outcomes[field] = {"status": "refused", "reason": "managed by your subscription/connection; change it in the Subscription section"} + elif write_would_suicide(field, value, powering): + outcomes[field] = {"status": "refused", "reason": f"would disconnect {powering.label}, which is powering this run"} + else: + staged[field] = value + + if staged: + merged = settings.model_dump() + merged.update(staged) + try: new_body = AppSettings(**merged) - if staged and new_body is not None: - await update_settings(new_body) - for f in staged: - outcomes[f] = {"status": "applied"} + except ValidationError as e: + bad = {str(err["loc"][0]) for err in e.errors() if err.get("loc")} + for f in bad & set(staged.keys()): + outcomes[f] = {"status": "refused", "reason": "invalid value for this field"} + staged.pop(f, None) + new_body = None + if staged: + merged = settings.model_dump() + merged.update(staged) + new_body = AppSettings(**merged) + if staged and new_body is not None: + try: + await update_settings(new_body) + for f in staged: + outcomes[f] = {"status": "applied"} + except Exception as e: + # Don't hand the agent an opaque 500; tell it which writes failed. + for f in staged: + outcomes[f] = {"status": "error", "reason": f"write failed: {e}"} return JSONResponse({"outcomes": outcomes}) diff --git a/backend/tests/test_settings_meta_concurrency.py b/backend/tests/test_settings_meta_concurrency.py new file mode 100644 index 00000000..eebdbe5e --- /dev/null +++ b/backend/tests/test_settings_meta_concurrency.py @@ -0,0 +1,62 @@ +"""Concurrent SettingsWrite must not lose updates. + +SettingsWrite is a read-modify-write that routes through update_settings (which +awaits), so two autonomous agents writing at the same time would interleave: each +loads the same snapshot, each writes the WHOLE object back, and the last writer +silently reverts the other's field, while BOTH agents are told "applied". This +drives two genuinely concurrent writes (different fields) through the ASGI app +and asserts neither is lost. It's the regression guard for the asyncio lock that +serializes these writes; without the lock this fails reproducibly. +""" + +from __future__ import annotations + +import asyncio + +import httpx +import pytest + +from backend.main import app + + +def _auth_headers(): + import backend.auth as auth_mod + if not auth_mod._TOKEN: + import secrets + auth_mod._TOKEN = secrets.token_urlsafe(32) + return {"Authorization": f"Bearer {auth_mod._TOKEN}"} + + +@pytest.fixture +def reset_settings(): + from backend.apps.settings.settings import load_settings, _save_settings + original = load_settings().model_copy(deep=True) + yield + _save_settings(original) + + +@pytest.mark.asyncio +async def test_concurrent_writes_to_different_fields_both_survive(reset_settings): + from backend.apps.settings.settings import load_settings, _save_settings + + base = load_settings() + base.theme = "dark" + base.default_mode = "agent" + _save_settings(base) + + headers = _auth_headers() + transport = httpx.ASGITransport(app=app) + async with httpx.AsyncClient(transport=transport, base_url="http://test", headers=headers) as client: + r1, r2 = await asyncio.gather( + client.post("/api/settings-meta/write", json={"changes": {"theme": "light"}}), + client.post("/api/settings-meta/write", json={"changes": {"default_mode": "chat"}}), + ) + + assert r1.status_code == 200 and r2.status_code == 200 + assert r1.json()["outcomes"]["theme"]["status"] == "applied" + assert r2.json()["outcomes"]["default_mode"]["status"] == "applied" + + final = load_settings() + # Both concurrent edits must persist; neither agent's "applied" result is a lie. + assert final.theme == "light", "lost update: theme was clobbered by the concurrent write" + assert final.default_mode == "chat", "lost update: default_mode was clobbered" diff --git a/backend/tests/test_skill_registry_community.py b/backend/tests/test_skill_registry_community.py index c910d75b..3ca1c713 100644 --- a/backend/tests/test_skill_registry_community.py +++ b/backend/tests/test_skill_registry_community.py @@ -84,6 +84,21 @@ def test_write_folder_skill_lands_files_and_indexes(skills_dir): assert "pdf-tk" in {s.id for s in skills_mod._sync_skills()} +def test_install_dedups_instead_of_clobbering_existing_skill(skills_dir): + # A user already has a local skill named "pdf". + skills_mod.write_folder_skill("pdf", {"SKILL.md": "MINE"}, {"name": "My PDF"}) + # A wild-registry install of a same-named skill must NOT overwrite it. + slug = skills_mod.unique_skill_slug("pdf") + assert slug == "pdf-2" + skills_mod.write_folder_skill(slug, {"SKILL.md": "THEIRS"}, {"name": "Registry PDF"}) + with open(skills_dir / "pdf" / "SKILL.md", encoding="utf-8") as f: + assert f.read() == "MINE", "registry install clobbered the user's existing skill" + with open(skills_dir / "pdf-2" / "SKILL.md", encoding="utf-8") as f: + assert f.read() == "THEIRS" + ids = {s.id for s in skills_mod._sync_skills()} + assert {"pdf", "pdf-2"} <= ids + + def test_write_folder_skill_blocks_path_traversal(skills_dir): skills_mod.write_folder_skill( "evil", diff --git a/frontend/src/app/pages/Skills/CommunitySkillsDialog.tsx b/frontend/src/app/pages/Skills/CommunitySkillsDialog.tsx index e5ed7a9e..8cf0ad1a 100644 --- a/frontend/src/app/pages/Skills/CommunitySkillsDialog.tsx +++ b/frontend/src/app/pages/Skills/CommunitySkillsDialog.tsx @@ -41,17 +41,25 @@ const CommunitySkillsDialog: React.FC = ({ open, onClose, onInstalled }) const [busy, setBusy] = useState(false); const [error, setError] = useState(null); const debounceRef = useRef | null>(null); + // Monotonic request tokens: a slow response from an earlier search/preview + // must not overwrite the state a newer one already set (out-of-order network). + const searchSeq = useRef(0); + const previewSeq = useRef(0); const runSearch = useCallback(async (q: string) => { + const seq = ++searchSeq.current; setLoading(true); setError(null); try { - setResults(await searchCommunitySkills(q)); + const res = await searchCommunitySkills(q); + if (seq !== searchSeq.current) return; + setResults(res); } catch (e) { + if (seq !== searchSeq.current) return; setError(e instanceof Error ? e.message : 'Search failed'); setResults([]); } finally { - setLoading(false); + if (seq === searchSeq.current) setLoading(false); } }, []); @@ -69,18 +77,21 @@ const CommunitySkillsDialog: React.FC = ({ open, onClose, onInstalled }) }, [open]); const preview = async (skill: CommunitySkill) => { + const seq = ++previewSeq.current; setSelected(skill); setDisclosure(null); setBusy(true); setError(null); try { const res = await installCommunitySkill(skill.source, skill.skillId, false); + if (seq !== previewSeq.current) return; setDisclosure(res.disclosure); } catch (e) { + if (seq !== previewSeq.current) return; setError(e instanceof Error ? e.message : 'Could not load skill'); setSelected(null); } finally { - setBusy(false); + if (seq === previewSeq.current) setBusy(false); } }; @@ -149,7 +160,7 @@ const CommunitySkillsDialog: React.FC = ({ open, onClose, onInstalled }) {selected && ( -