diff --git a/lib/galaxy/tools/cached_toolbox.py b/lib/galaxy/tools/cached_toolbox.py index 026683ce77b..d4a2c4e2bd9 100644 --- a/lib/galaxy/tools/cached_toolbox.py +++ b/lib/galaxy/tools/cached_toolbox.py @@ -846,7 +846,7 @@ class CachedToolBox(ToolBox): load instead of degrading silently. """ log.info("CachedToolBox: running populator inline to backfill the index") - populate_store_inline(self.app.config) + populate_store_inline(self.app.config, app=self.app) def _index_versions_for(self, tool_id: str) -> list[str]: """Return every version present in the index for ``tool_id``. diff --git a/lib/galaxy/tools/source_store/populator.py b/lib/galaxy/tools/source_store/populator.py index 48df17c3b75..e652f3968eb 100644 --- a/lib/galaxy/tools/source_store/populator.py +++ b/lib/galaxy/tools/source_store/populator.py @@ -607,7 +607,7 @@ def populate_store( ) -def populate_store_inline( +def _populate_store_inline_unlocked( config: GalaxyAppConfiguration, *, paths: list[str] | None = None, @@ -984,6 +984,54 @@ def populate_store_inline( return stats +def populate_store_inline( + config: GalaxyAppConfiguration, + *, + paths: list[str] | None = None, + pattern: str | None = None, + parallel: int = 1, + dry_run: bool = False, + incremental: bool = True, + verbose: bool = False, + broadcast: bool = False, + target: str | None = None, + prune: bool = False, + path_guids: dict[str, str | None] | None = None, + app=None, + write_manifests: bool = False, +) -> dict[str, int]: + """Populate a store, serialized with toolbox mutations in this process. + + Discovery and the eventual index write are one transaction with respect to + in-process installs, uninstalls, and toolbox replacement. Without this lock + a populator can discover a shed tool before its conf entry is removed, then + commit that stale snapshot after ``remove_index_entry`` and resurrect the + uninstalled tool. + """ + + def populate() -> dict[str, int]: + return _populate_store_inline_unlocked( + config, + paths=paths, + pattern=pattern, + parallel=parallel, + dry_run=dry_run, + incremental=incremental, + verbose=verbose, + broadcast=broadcast, + target=target, + prune=prune, + path_guids=path_guids, + app=app, + write_manifests=write_manifests, + ) + + if app is None: + return populate() + with app._toolbox_lock: + return populate() + + def populate_for_paths( config: GalaxyAppConfiguration, paths: list[str], diff --git a/test/unit/scripts/tool_source/test_populate_store.py b/test/unit/scripts/tool_source/test_populate_store.py index e78c3d8f4cc..be8d55467dd 100644 --- a/test/unit/scripts/tool_source/test_populate_store.py +++ b/test/unit/scripts/tool_source/test_populate_store.py @@ -10,8 +10,10 @@ Following test guidelines: """ import json +import threading from dataclasses import dataclass from pathlib import Path +from types import SimpleNamespace from typing import ( Any, ) @@ -22,6 +24,7 @@ import sys sys.path.insert(0, str(galaxy_root / "lib")) +from galaxy.tools.source_store.discover import discover_tools as actual_discover_tools from galaxy.tools.source_store.factory import _build_default_store from galaxy.tools.source_store.freshness import tool_confs_token from galaxy.tools.source_store.populator import ( @@ -411,6 +414,67 @@ class TestIncrementalFastPath: assert entry.panel_section_id == "test_section_multi" assert [(item.tool_id, item.section_id) for item in index.panel_items] == [(guid, "test_section_multi")] + def test_in_process_population_is_atomic_with_tool_removal(self, tmp_path, monkeypatch): + """A stale scan must commit before a concurrent uninstall can prune it.""" + tools_dir = tmp_path / "tools" + _write_tool(tools_dir / "keep.xml", "keep") + _write_tool(tools_dir / "doomed.xml", "doomed") + conf = _write_tool_conf( + tmp_path / "tool_conf.xml", + tools_dir, + '', + ) + cfg = _populate_config(tmp_path, conf) + _populate_default(cfg) + + from galaxy.tools.source_store import populator + + scan_complete = threading.Event() + allow_commit = threading.Event() + + def pause_after_discovery(*args, **kwargs): + discovered = list(actual_discover_tools(*args, **kwargs)) + scan_complete.set() + if not allow_commit.wait(timeout=5): + raise TimeoutError("test did not release the paused population") + return iter(discovered) + + monkeypatch.setattr(populator, "discover_tools", pause_after_discovery) + app = SimpleNamespace(_toolbox_lock=threading.RLock()) + errors: list[Exception] = [] + + def populate_from_stale_scan(): + try: + _populate_default(cfg, app=app, incremental=False) + except Exception as exc: + errors.append(exc) + + population = threading.Thread(target=populate_from_stale_scan) + population.start() + assert scan_complete.wait(timeout=5) + + # With no population lock this succeeds: uninstall prunes the index, + # then the paused stale scan writes ``doomed`` straight back. + acquired_during_population = app._toolbox_lock.acquire(blocking=False) + if acquired_during_population: + _write_tool_conf(conf, tools_dir, '') + _build_default_store(cfg).remove_index_entry("doomed") + app._toolbox_lock.release() + + allow_commit.set() + population.join(timeout=10) + assert not population.is_alive() + assert not errors + + if not acquired_during_population: + with app._toolbox_lock: + _write_tool_conf(conf, tools_dir, '') + _build_default_store(cfg).remove_index_entry("doomed") + + assert not acquired_during_population + _, index = _load_default_index(cfg) + assert "doomed" not in index.entries + def test_manifest_is_opt_in_for_cli_callers(self, tmp_path): tools_dir = tmp_path / "tools" _write_tool(tools_dir / "manifest_tool.xml", "manifest_tool", name="Manifest")