mirror of
https://github.com/openswarm-ai/openswarm.git
synced 2026-10-01 14:04:51 +02:00
[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
This commit is contained in:
@@ -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)
|
||||
|
||||
@@ -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"]},
|
||||
)
|
||||
|
||||
@@ -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
|
||||
|
||||
+51
-35
@@ -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})
|
||||
|
||||
|
||||
@@ -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"
|
||||
@@ -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",
|
||||
|
||||
@@ -41,17 +41,25 @@ const CommunitySkillsDialog: React.FC<Props> = ({ open, onClose, onInstalled })
|
||||
const [busy, setBusy] = useState(false);
|
||||
const [error, setError] = useState<string | null>(null);
|
||||
const debounceRef = useRef<ReturnType<typeof setTimeout> | 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<Props> = ({ 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<Props> = ({ open, onClose, onInstalled })
|
||||
|
||||
{selected && (
|
||||
<Box sx={{ display: 'flex', flexDirection: 'column', gap: 1 }}>
|
||||
<Button onClick={() => { setSelected(null); setDisclosure(null); }} size="small"
|
||||
<Button onClick={() => { previewSeq.current++; setSelected(null); setDisclosure(null); }} size="small"
|
||||
sx={{ alignSelf: 'flex-start', textTransform: 'none', color: c.text.tertiary, fontSize: '0.78rem' }}>
|
||||
← Back to results
|
||||
</Button>
|
||||
|
||||
Reference in New Issue
Block a user