Merge & dedup data table entries on shed install (no DM gate)

Replace the non-Data-Manager skip with a uniform install path: every
repo's tables get registered. First install copies the .loc.sample to
tool_data_path and writes the <table> entry without a
<tool_shed_repository> sub-element; subsequent installs of the same
table append unique rows from their .loc.sample to the shared loc file,
prefixed by a comment_char attribution line. get_filename_for_source
falls back to a no-repo_info filename so DMs on legacy installs keep
matching their existing <tool_shed_repository>-tagged files.

Co-Authored-By: Claude Opus 4.7 <noreply@anthropic.com>
This commit is contained in:
Nate Coraor
2026-05-14 14:07:57 -04:00
co-authored by Claude Opus 4.7
parent 07c5d445da
commit 3e00f529ba
8 changed files with 237 additions and 21 deletions
@@ -171,7 +171,7 @@ class InstallRepositoryManager:
session.commit()
is_data_manager = "data_manager" in irmm_metadata_dict
if is_data_manager and "sample_files" in irmm_metadata_dict:
if "sample_files" in irmm_metadata_dict:
sample_files = irmm_metadata_dict.get("sample_files", [])
tool_index_sample_files = stdtm.get_tool_index_sample_files(sample_files)
tool_data_table_conf_filename, tool_data_table_elems = stdtm.install_tool_data_tables(
@@ -8,7 +8,10 @@ from typing import (
from galaxy.tool_shed.galaxy_install.client import InstallationTarget
from galaxy.tool_shed.util import hg_util
from galaxy.tool_util.data import DataTableColumnMismatch
from galaxy.tool_util.data import (
DataTableColumnMismatch,
TabularToolDataTable,
)
from galaxy.util import (
Element,
SubElement,
@@ -28,8 +31,6 @@ RequiredAppT = Union["BasicSharedApp", InstallationTarget]
def _parse_table_columns(table_elem: Element) -> dict[str, int]:
"""Parse a ``<table>`` element's column spec into a name->index mapping."""
from galaxy.tool_util.data import TabularToolDataTable
columns, _, _ = TabularToolDataTable.parse_column_spec_element(table_elem)
if "value" in columns and "name" not in columns:
columns["name"] = columns["value"]
@@ -157,6 +158,28 @@ class ShedToolDataTableManager(BaseShedToolDataTableManager):
os.makedirs(target_dir)
return target_dir, tool_path, relative_target_dir
def _merge_loc_sample_entries(
self,
existing_table: TabularToolDataTable,
elem: Element,
loc_basename_to_source: dict,
tool_shed_repository: "ToolShedRepository",
) -> None:
attribution = (
f"Added by {tool_shed_repository.owner}/{tool_shed_repository.name}"
f"@{tool_shed_repository.installed_changeset_revision}"
)
for file_elem in elem.findall("file"):
shared_path = file_elem.get("path")
basename = os.path.basename(shared_path) if shared_path else None
source_loc_sample = loc_basename_to_source.get(basename) if basename else None
if not source_loc_sample or not os.path.exists(source_loc_sample):
continue
# Keep `${__HERE__}` literal so the appended rows match the shared loc's existing format.
new_rows = existing_table.parse_file_fields(source_loc_sample, here="${__HERE__}")
if new_rows:
existing_table.append_entries_with_attribution(new_rows, attribution)
def install_tool_data_tables(self, tool_shed_repository: "ToolShedRepository", tool_index_sample_files):
TOOL_DATA_TABLE_FILE_NAME = "tool_data_table_conf.xml"
TOOL_DATA_TABLE_FILE_SAMPLE_NAME = f"{TOOL_DATA_TABLE_FILE_NAME}.sample"
@@ -164,8 +187,10 @@ class ShedToolDataTableManager(BaseShedToolDataTableManager):
SAMPLE_SUFFIX_OFFSET = -len(SAMPLE_SUFFIX)
LOC_SAMPLE_SUFFIX = ".loc.sample"
target_dir, tool_path, relative_target_dir = self.get_target_install_dir(tool_shed_repository)
# .loc files are shared across DM revisions: install at the root of tool_data_path.
# .loc files are shared across repository installations: install at the root of tool_data_path.
shared_loc_dir = self.app.config.tool_data_path
# Map shared loc basename -> source .loc.sample, used to merge entries when a table is reinstalled.
loc_basename_to_source: dict[str, str] = {}
for sample_file in tool_index_sample_files:
path, filename = os.path.split(sample_file)
target_filename = filename
@@ -175,6 +200,7 @@ class ShedToolDataTableManager(BaseShedToolDataTableManager):
if filename.endswith(LOC_SAMPLE_SUFFIX):
target_path_filename = os.path.join(shared_loc_dir, target_filename)
install_dest_dir = shared_loc_dir
loc_basename_to_source[target_filename] = source_file
else:
target_path_filename = os.path.join(target_dir, target_filename)
install_dest_dir = target_dir
@@ -225,11 +251,11 @@ class ShedToolDataTableManager(BaseShedToolDataTableManager):
incoming_columns = _parse_table_columns(elem)
self.app.tool_data_tables.assert_data_table_consistency(table_name, incoming_columns)
existing = registered_tables.get(table_name) if table_name else None
if existing is not None and getattr(existing, "columns", None) is not None:
# Already registered with matching columns; skip to avoid duplicate <table> entry.
if isinstance(existing, TabularToolDataTable) and existing.columns is not None:
# Already registered with matching columns. Merge any rows from this install's
# .loc.sample(s) into the shared loc file, then skip adding a duplicate <table>.
self._merge_loc_sample_entries(existing, elem, loc_basename_to_source, tool_shed_repository)
continue
# Store repository info in the table tag set for trace-ability.
self.generate_repository_info_elem_from_repository(tool_shed_repository, parent_elem=elem)
kept_elems.append(elem)
if kept_elems:
# Remove old data_table
+65
View File
@@ -684,6 +684,7 @@ class TabularToolDataTable(ToolDataTable):
source_repo_info_model = source.repo_info
source_repo_info = source_repo_info_model.model_dump() if source_repo_info_model else None
filename = default
shared_fallback: Optional[str] = None
for name, value in self.filenames.items():
repo_info = value.get("tool_shed_repository")
if (not source_repo_info and not repo_info) or (
@@ -691,6 +692,14 @@ class TabularToolDataTable(ToolDataTable):
):
filename = name
break
if source_repo_info and not repo_info and shared_fallback is None:
# Loc file registered without repo_info (shared install) — use as fallback
# when a tool-shed source has no exact match. Preserves legacy behavior where
# a `<tool_shed_repository>`-tagged filename matches exactly.
shared_fallback = name
else:
if shared_fallback is not None:
filename = shared_fallback
return filename
def _add_entry(
@@ -762,6 +771,62 @@ class TabularToolDataTable(ToolDataTable):
fields_collapsed = f"{self.separator.join(fields)}\n"
data_table_fh.write(fields_collapsed.encode("utf-8"))
def append_entries_with_attribution(
self,
entries: List[List[str]],
attribution: str,
allow_duplicates: bool = False,
) -> int:
"""Append ``entries`` to this table's loc file with an attribution comment line.
When ``allow_duplicates`` is False, dedupes entries by the ``value`` column
(the table's primary key) against existing in-memory rows and the pending
batch. Writes a single ``{comment_char} {attribution}`` line before the first
new row. No-op if no rows survive the dedup.
"""
filename: Optional[str] = self.get_filename_for_source(None)
if filename is None:
for name in self.filenames:
filename = name
break
if filename is None:
raise MessageException(f"Unable to determine filename for appending entries to data table '{self.name}'.")
value_index = self.columns.get("value", 0)
existing_values: Optional[Set[str]] = None
if not allow_duplicates:
existing_values = {row[value_index] for row in self.data if value_index < len(row)}
new_rows: List[List[str]] = []
for entry in entries:
fields = self._replace_field_separators(list(entry))
if self.largest_index >= len(fields):
log.warning(
"Skipping fields (%s) for data table '%s': only %d field(s) provided.",
fields,
self.name,
len(fields),
)
continue
if existing_values is not None and fields[value_index] in existing_values:
continue
new_rows.append(fields)
if existing_values is not None:
existing_values.add(fields[value_index])
if not new_rows:
return self._loaded_content_version
with FileLock(filename):
if os.path.exists(filename) and os.stat(filename).st_size > 0:
with open(filename, "rb+") as fh:
fh.seek(-1, 2)
if fh.read(1) not in (b"\n", b"\r"):
fh.write(b"\n")
with open(filename, "ab") as fh:
fh.write(f"{self.comment_char} {attribution}\n".encode())
for fields in new_rows:
fh.write(f"{self.separator.join(fields)}\n".encode())
for fields in new_rows:
self.data.append(fields)
return self._update_version()
def _locate_filename(self, tool_data_file_path, entry_source, bundle_mode):
if tool_data_file_path is not None:
filename = tool_data_file_path
+17 -5
View File
@@ -62,15 +62,27 @@ class TestRepositoryInstallIntegrationTestCase(integration_util.IntegrationTestC
self.install_repository(*repo)
self.uninstall_repository(*repo)
def test_non_data_manager_install_skips_data_table_registration(self):
"""Non-Data-Manager repos must not persist data table entries in shed_tool_data_table_conf.xml."""
def test_non_data_manager_install_registers_data_tables_without_repo_info(self):
"""Non-Data-Manager repos register tables in shed_tool_data_table_conf.xml without the
<tool_shed_repository> sub-element, and their loc files land at tool_data_path root."""
non_dm_repo = ("devteam", "bwa", "051eba708f43")
non_dm_table_names = {"bwa_indexes", "bwa_mem_indexes"}
self.install_repository(*non_dm_repo)
shed_conf = self._app.config.shed_tool_data_table_config
registered = {t.get("name") for t in ET.parse(shed_conf).getroot().findall("table")}
leaked = non_dm_table_names & registered
assert not leaked, f"Unexpected data tables in {shed_conf}: {sorted(leaked)}"
table_elems = {t.get("name"): t for t in ET.parse(shed_conf).getroot().findall("table")}
missing = non_dm_table_names - table_elems.keys()
assert not missing, f"Expected tables not registered in {shed_conf}: {sorted(missing)}"
for name in non_dm_table_names:
assert (
table_elems[name].find("tool_shed_repository") is None
), f"Table {name!r} should not have a <tool_shed_repository> sub-element on new installs"
file_elems = table_elems[name].findall("file")
assert file_elems, f"Table {name!r} has no <file> entry"
for file_elem in file_elems:
loc_path = file_elem.get("path") or ""
assert loc_path.startswith(
self._app.config.tool_data_path
), f"Loc file for {name!r} not at tool_data_path root: {loc_path}"
def test_repository_update(self):
response = self._install_repository(revision=REVISION_4, version="0.0.3", allow_upgraded=True)[0]
+3 -2
View File
@@ -43,7 +43,8 @@ def test_against_production_shed(tmp_path: Path):
assert tool_guid in f.read()
repo_path = tmp_path / "tools" / "toolshed.g2.bx.psu.edu" / "repos" / repo_owner / repo_name / repo_revision
assert repo_path.exists()
# featurecounts is not a Data Manager — install must not register a per-revision data table config.
# All shed installs now register data tables; the per-revision config is written for any repo
# with sample files, regardless of whether it's a Data Manager.
tool_data_table_path = (
tmp_path
/ "tool_data"
@@ -54,7 +55,7 @@ def test_against_production_shed(tmp_path: Path):
/ repo_revision
/ "tool_data_table_conf.xml"
)
assert not tool_data_table_path.exists()
assert tool_data_table_path.exists()
install_model_context = cast("install_model_scoped_session", install_target.install_model.session)
query = install_model_context.query(ToolShedRepository).where(ToolShedRepository.name == repo_name)
+44 -2
View File
@@ -9,7 +9,10 @@ from galaxy.tool_shed.tools.data_table_manager import (
DataTableColumnMismatch,
ShedToolDataTableManager,
)
from galaxy.tool_util.data import ToolDataTableManager
from galaxy.tool_util.data import (
TabularToolDataTable,
ToolDataTableManager,
)
from galaxy.util import (
Element,
SubElement,
@@ -80,9 +83,10 @@ def _make_stdtm(tmp_path):
def _registered_table(columns, filenames=None):
existing = mock.MagicMock(spec=["columns", "filenames"])
existing = mock.MagicMock(spec=TabularToolDataTable)
existing.columns = columns
existing.filenames = filenames or {}
existing.parse_file_fields.return_value = []
return existing
@@ -99,6 +103,10 @@ def test_loc_file_lands_at_shared_root_not_per_revision(tmp_path):
file_elems = list(kept_elems[0].findall("file"))
assert len(file_elems) == 1
assert file_elems[0].get("path") == shared_loc
# New installs should not stamp a <tool_shed_repository> sub-element on the <table>:
# the loc-file location is deterministic and DMs on legacy installs fall through to
# the shared-no-repo_info match in get_filename_for_source.
assert kept_elems[0].find("tool_shed_repository") is None
def test_existing_loc_file_is_not_overwritten(tmp_path):
@@ -163,6 +171,40 @@ def test_column_match_with_column_elements_dedupes(tmp_path):
assert kept_elems == []
def test_second_install_merges_loc_sample_rows_with_attribution(tmp_path):
stdtm, repo, samples, captured, _, _, _ = _make_stdtm(tmp_path)
matching_columns = {"value": 0, "dbkey": 1, "name": 2, "path": 3}
existing = _registered_table(matching_columns)
incoming_rows = [["hg19", "hg19", "human (hg19)", "/data/hg19.fa"]]
existing.parse_file_fields.return_value = incoming_rows
stdtm.app.tool_data_tables.data_tables = {"all_fasta": existing}
_, kept_elems = stdtm.install_tool_data_tables(repo, samples)
assert kept_elems == []
assert not captured["to_xml_calls"]
existing.append_entries_with_attribution.assert_called_once()
call_args = existing.append_entries_with_attribution.call_args
assert call_args.args[0] == incoming_rows
attribution = call_args.args[1]
assert "iuc/data_manager_fetch_genome_dbkeys_all_fasta" in attribution
assert "abc" in attribution
def test_second_install_with_empty_loc_sample_does_not_append(tmp_path):
stdtm, repo, samples, captured, _, _, _ = _make_stdtm(tmp_path)
matching_columns = {"value": 0, "dbkey": 1, "name": 2, "path": 3}
existing = _registered_table(matching_columns)
existing.parse_file_fields.return_value = []
stdtm.app.tool_data_tables.data_tables = {"all_fasta": existing}
_, kept_elems = stdtm.install_tool_data_tables(repo, samples)
assert kept_elems == []
assert not captured["to_xml_calls"]
existing.append_entries_with_attribution.assert_not_called()
def test_parse_table_columns_aliases_name_to_value():
from galaxy.tool_shed.tools.data_table_manager import _parse_table_columns
+3 -3
View File
@@ -69,11 +69,11 @@ def _invoke_handle(irm, metadata_dict: dict[str, Any], repository_tools_tups: li
return stdtm_instance
def test_non_data_manager_repo_skips_sample_files_registration():
def test_non_data_manager_repo_registers_sample_files():
irm, register_mock = _make_install_repository_manager()
stdtm_instance = _invoke_handle(irm, {"sample_files": ["foo.loc.sample"]}, [])
stdtm_instance.install_tool_data_tables.assert_not_called()
register_mock.assert_not_called()
stdtm_instance.install_tool_data_tables.assert_called_once()
register_mock.assert_called_once()
def test_data_manager_repo_registers_sample_files():
@@ -93,3 +93,73 @@ def test_assert_data_table_consistency_raises_column_mismatch(tdt_manager):
{"value": 0, "name": 1, "path": 2, "extra": 3},
)
assert exc_info.value.table_name == "testalpha"
def test_get_filename_for_source_falls_back_to_shared_filename(tdt_manager):
table = tdt_manager["testalpha"]
[shared_filename] = list(table.filenames)
assert table.filenames[shared_filename].get("tool_shed_repository") is None
source_with_unknown_repo = {
"tool_shed": "tool-shed",
"repository_name": "repo",
"repository_owner": "owner",
"installed_changeset_revision": "abc",
}
assert table.get_filename_for_source(source_with_unknown_repo) == shared_filename
def test_get_filename_for_source_prefers_exact_repo_match_over_shared(tdt_manager):
table = tdt_manager["testalpha"]
[shared_filename] = list(table.filenames)
legacy_info = {
"tool_shed": "tool-shed",
"repository_name": "legacy",
"repository_owner": "owner",
"installed_changeset_revision": "deadbeef",
}
legacy_filename = f"{shared_filename}.legacy"
table.filenames[legacy_filename] = dict(
found=True,
filename=legacy_filename,
from_shed_config=True,
tool_data_path=None,
config_element=None,
tool_shed_repository=legacy_info,
errors=[],
)
assert table.get_filename_for_source(legacy_info) == legacy_filename
other_info = dict(legacy_info, repository_name="unknown")
assert table.get_filename_for_source(other_info) == shared_filename
def test_append_entries_with_attribution_appends_and_dedupes(tdt_manager):
table = tdt_manager["testalpha"]
[loc_filename] = list(table.filenames)
initial_rows = len(table.data)
new_entries = [
["data3", "data3name", "${__HERE__}/data3/entry.txt"],
["data1", "data1name", "${__HERE__}/data1/entry.txt"], # duplicate, must be skipped
]
table.append_entries_with_attribution(new_entries, "added by owner/foo@rev1")
with open(loc_filename) as fh:
contents = fh.read()
assert "# added by owner/foo@rev1" in contents
assert contents.count("data3\tdata3name") == 1
assert contents.count("data1\tdata1name") == 1
assert len(table.data) == initial_rows + 1
def test_append_entries_with_attribution_noop_when_all_duplicates(tdt_manager):
table = tdt_manager["testalpha"]
[loc_filename] = list(table.filenames)
with open(loc_filename) as fh:
before = fh.read()
rows_before = list(table.data)
table.append_entries_with_attribution(
[["data1", "data1name", "${__HERE__}/data1/entry.txt"]],
"added by owner/foo@rev1",
)
with open(loc_filename) as fh:
after = fh.read()
assert after == before
assert table.data == rows_before