mirror of
https://github.com/galaxyproject/galaxy.git
synced 2026-09-24 16:30:27 +08:00
DatabaseToolSourceStore: own tool_source_record table, stop sharing tool_source
The job-request path creates and reads tool_source rows whose source column carries the raw source string; the store was writing dict-shaped payloads into the same table, and once identity hashes aligned the job path could resolve a store row and fail pydantic validation (500 on every request-style tool execution under lazy mode). Give the store its own content-addressed table with real columns — get_by_tool_id / get_by_source_path become indexed queries instead of full-table JSON scans, populator pruning can never delete rows the job path references, and the identity_hash coupling (and the legacy __tool_index__ sentinel-row fallback) disappear.
This commit is contained in:
@@ -53,15 +53,18 @@ Two persistence concepts:
|
||||
**StoredToolSource** — the canonical macro-expanded XML/YAML for a tool,
|
||||
keyed by SHA-256 of the expanded content. Multiple versions of the same
|
||||
``tool_id`` coexist as separate hashes. The DB backend persists these in the
|
||||
``tool_source`` table; the Redis and disk backends use their own layout.
|
||||
store-owned ``tool_source_record`` table (the ``tool_source`` table belongs
|
||||
to the job-request path and has a different payload contract); the Redis and
|
||||
disk backends use their own layout.
|
||||
|
||||
**ToolIndex** — a single dataclass containing one ``ToolIndexEntry`` per tool,
|
||||
holding everything the batch APIs need (id, name, description, panel section,
|
||||
labels, EDAM, requirements, container info, test counts, hidden/disabled,
|
||||
shed metadata). The index is serialized and gzip-compressed as a blob.
|
||||
|
||||
The DB backend gets a new ``tool_index`` table (migration
|
||||
``f5a73c8b9d12_add_tool_index_table``) with a single row per index version.
|
||||
The DB backend gets the ``tool_index`` and ``tool_source_record`` tables
|
||||
(migration ``f5a73c8b9d12``); ``tool_index`` holds a single row per index
|
||||
version.
|
||||
Redis stores the blob under a known key; disk stores it as a file.
|
||||
|
||||
Backend Abstraction
|
||||
|
||||
@@ -67,19 +67,11 @@ def get_or_create_tool_source(session: Session, tool) -> ToolSource:
|
||||
|
||||
def tool_source_identity_hash(tool: Any) -> str:
|
||||
dynamic_tool = getattr(tool, "dynamic_tool", None)
|
||||
identity: tuple[str, ...]
|
||||
if dynamic_tool is not None and dynamic_tool.id is not None:
|
||||
identity = ("dynamic", str(dynamic_tool.id))
|
||||
return hashlib.sha256("\0".join(identity).encode("utf-8")).hexdigest()
|
||||
return static_tool_source_identity_hash(tool.id, tool.version)
|
||||
|
||||
|
||||
def static_tool_source_identity_hash(tool_id: str | None, tool_version: str | None) -> str:
|
||||
"""Identity hash for a persisted source with no dynamic-tool linkage.
|
||||
|
||||
Shared with ``galaxy.tool_source_store`` writers, which persist
|
||||
config-file tool sources and therefore always carry a static identity.
|
||||
"""
|
||||
identity = ("static", tool_id or "", tool_version or "")
|
||||
else:
|
||||
identity = ("static", tool.id or "", tool.version or "")
|
||||
return hashlib.sha256("\0".join(identity).encode("utf-8")).hexdigest()
|
||||
|
||||
|
||||
|
||||
@@ -1428,6 +1428,36 @@ class ToolSource(Base, Dictifiable, RepresentById):
|
||||
dynamic_tool: Mapped[Optional["DynamicTool"]] = relationship()
|
||||
|
||||
|
||||
class ToolSourceRecord(Base, Dictifiable, RepresentById):
|
||||
"""
|
||||
A tool source persisted by the tool source store.
|
||||
|
||||
Deliberately separate from :class:`ToolSource`, whose rows are created
|
||||
per executed tool by the job-request path and whose ``source`` column
|
||||
carries the raw source string contract that path deserializes. Store
|
||||
rows are content-addressed by ``hash`` and carry populator metadata;
|
||||
the populator may prune them freely without affecting job records.
|
||||
"""
|
||||
|
||||
__tablename__ = "tool_source_record"
|
||||
|
||||
dict_collection_visible_keys = ("id", "hash", "tool_id", "tool_version", "create_time", "update_time")
|
||||
dict_element_visible_keys = ("id", "hash", "tool_id", "tool_version", "create_time", "update_time")
|
||||
|
||||
id: Mapped[int] = mapped_column(primary_key=True)
|
||||
hash: Mapped[str] = mapped_column(String(255), unique=True, index=True, nullable=False)
|
||||
source: Mapped[str] = mapped_column(Text().with_variant(mysql.LONGTEXT(), "mysql"), nullable=False)
|
||||
source_class: Mapped[str] = mapped_column(TrimmedString(255))
|
||||
tool_id: Mapped[str | None] = mapped_column(String(255), index=True)
|
||||
tool_version: Mapped[str | None] = mapped_column(String(255))
|
||||
tool_dir: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
source_path: Mapped[str | None] = mapped_column(Text, nullable=True)
|
||||
stored_at: Mapped[datetime | None] = mapped_column(DateTime, nullable=True)
|
||||
source_metadata: Mapped[dict | None] = mapped_column(JSONType, nullable=True)
|
||||
create_time: Mapped[datetime] = mapped_column(DateTime, default=now, nullable=False)
|
||||
update_time: Mapped[datetime] = mapped_column(DateTime, default=now, onupdate=now, nullable=False)
|
||||
|
||||
|
||||
class ToolIndexCache(Base, Dictifiable, RepresentById):
|
||||
"""
|
||||
Stores pre-computed tool index for fast API responses.
|
||||
|
||||
+34
-8
@@ -1,4 +1,4 @@
|
||||
"""Add tool_index table for storing pre-computed tool index
|
||||
"""Add tool source store tables (tool_index, tool_source_record)
|
||||
|
||||
Revision ID: f5a73c8b9d12
|
||||
Revises: 28885b317f78
|
||||
@@ -9,6 +9,7 @@ Create Date: 2026-01-25 10:00:00.000000
|
||||
import sqlalchemy as sa
|
||||
from sqlalchemy.dialects import mysql
|
||||
|
||||
from galaxy.model.custom_types import JSONType
|
||||
from galaxy.model.migrations.util import (
|
||||
create_table,
|
||||
drop_table,
|
||||
@@ -20,18 +21,21 @@ down_revision = "28885b317f78"
|
||||
branch_labels = None
|
||||
depends_on = None
|
||||
|
||||
TABLE_NAME = "tool_index"
|
||||
INDEX_TABLE_NAME = "tool_index"
|
||||
SOURCE_TABLE_NAME = "tool_source_record"
|
||||
|
||||
|
||||
def upgrade():
|
||||
"""Create tool_index table for storing pre-computed tool index.
|
||||
"""Create the tool source store's tables.
|
||||
|
||||
This table stores a serialized ToolIndex object that provides
|
||||
fast access to tool metadata for API responses without loading
|
||||
full tool sources.
|
||||
``tool_index`` stores a serialized ToolIndex object that provides fast
|
||||
access to tool metadata for API responses without loading full tool
|
||||
sources. ``tool_source_record`` stores the content-addressed tool
|
||||
sources themselves — separate from ``tool_source``, whose rows belong
|
||||
to the job-request path and carry a different payload contract.
|
||||
"""
|
||||
create_table(
|
||||
TABLE_NAME,
|
||||
INDEX_TABLE_NAME,
|
||||
sa.Column("id", sa.Integer, primary_key=True),
|
||||
sa.Column("version", sa.String(64), nullable=False, unique=True),
|
||||
sa.Column("data", sa.LargeBinary().with_variant(mysql.LONGBLOB(), "mysql"), nullable=False),
|
||||
@@ -45,7 +49,29 @@ def upgrade():
|
||||
onupdate=sa.func.now(),
|
||||
),
|
||||
)
|
||||
create_table(
|
||||
SOURCE_TABLE_NAME,
|
||||
sa.Column("id", sa.Integer, primary_key=True),
|
||||
sa.Column("hash", sa.String(255), nullable=False, unique=True, index=True),
|
||||
sa.Column("source", sa.Text().with_variant(mysql.LONGTEXT(), "mysql"), nullable=False),
|
||||
sa.Column("source_class", sa.String(255)),
|
||||
sa.Column("tool_id", sa.String(255), index=True),
|
||||
sa.Column("tool_version", sa.String(255)),
|
||||
sa.Column("tool_dir", sa.Text, nullable=True),
|
||||
sa.Column("source_path", sa.Text, nullable=True),
|
||||
sa.Column("stored_at", sa.DateTime, nullable=True),
|
||||
sa.Column("source_metadata", JSONType, nullable=True),
|
||||
sa.Column("create_time", sa.DateTime, nullable=False, server_default=sa.func.now()),
|
||||
sa.Column(
|
||||
"update_time",
|
||||
sa.DateTime,
|
||||
nullable=False,
|
||||
server_default=sa.func.now(),
|
||||
onupdate=sa.func.now(),
|
||||
),
|
||||
)
|
||||
|
||||
|
||||
def downgrade():
|
||||
drop_table(TABLE_NAME)
|
||||
drop_table(SOURCE_TABLE_NAME)
|
||||
drop_table(INDEX_TABLE_NAME)
|
||||
|
||||
@@ -2,8 +2,9 @@
|
||||
Database backend for Tool Source Store.
|
||||
|
||||
This module provides a database-backed implementation of the ToolSourceStore
|
||||
that uses the existing tool_source table and adds a tool_index table for
|
||||
lightweight metadata.
|
||||
backed by the store-owned ``tool_source_record`` table (content-addressed
|
||||
sources) and the ``tool_index`` table (serialized ToolIndex). The unrelated
|
||||
``tool_source`` table belongs to the job-request path and is never touched.
|
||||
"""
|
||||
|
||||
import gzip
|
||||
@@ -11,7 +12,6 @@ import json
|
||||
import logging
|
||||
from collections.abc import Iterator
|
||||
from contextlib import contextmanager
|
||||
from datetime import datetime
|
||||
from typing import (
|
||||
cast,
|
||||
)
|
||||
@@ -22,10 +22,9 @@ from sqlalchemy import (
|
||||
)
|
||||
from sqlalchemy.orm import Session
|
||||
|
||||
from galaxy.managers.tool_source import static_tool_source_identity_hash
|
||||
from galaxy.model import (
|
||||
ToolIndexCache,
|
||||
ToolSource as ToolSourceModel,
|
||||
ToolSourceRecord,
|
||||
)
|
||||
from galaxy.model.scoped_session import galaxy_scoped_session
|
||||
from . import (
|
||||
@@ -44,8 +43,8 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
"""
|
||||
Database-backed tool source store.
|
||||
|
||||
Uses the existing tool_source table for storing full tool sources
|
||||
and a separate tool_index table for lightweight metadata.
|
||||
Uses the ``tool_source_record`` table for full tool sources and the
|
||||
``tool_index`` table for the serialized index.
|
||||
"""
|
||||
|
||||
def __init__(self, sa_session: galaxy_scoped_session):
|
||||
@@ -89,35 +88,22 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
"""Store a tool source in the database."""
|
||||
session = self._get_session()
|
||||
|
||||
# Check if already exists
|
||||
existing = session.execute(
|
||||
select(ToolSourceModel).where(ToolSourceModel.hash == tool_source.hash)
|
||||
select(ToolSourceRecord.id).where(ToolSourceRecord.hash == tool_source.hash)
|
||||
).scalar_one_or_none()
|
||||
|
||||
if existing:
|
||||
return tool_source.hash
|
||||
|
||||
# Create new record
|
||||
source_data = {
|
||||
"raw": tool_source.raw_source,
|
||||
"tool_source_class": tool_source.tool_source_class,
|
||||
"tool_id": tool_source.tool_id,
|
||||
"tool_version": tool_source.tool_version,
|
||||
"tool_dir": tool_source.tool_dir,
|
||||
"source_path": tool_source.source_path,
|
||||
"stored_at": (tool_source.stored_at.isoformat() if tool_source.stored_at else None),
|
||||
"metadata": tool_source.metadata,
|
||||
}
|
||||
|
||||
model = ToolSourceModel(
|
||||
model = ToolSourceRecord(
|
||||
hash=tool_source.hash,
|
||||
source=source_data,
|
||||
source=tool_source.raw_source,
|
||||
source_class=tool_source.tool_source_class,
|
||||
tool_id=tool_source.tool_id,
|
||||
tool_version=tool_source.tool_version,
|
||||
# Same "static" identity as galaxy.managers.tool_source.tool_source_identity_hash —
|
||||
# store-backed sources never carry a dynamic tool.
|
||||
identity_hash=static_tool_source_identity_hash(tool_source.tool_id, tool_source.tool_version),
|
||||
tool_dir=tool_source.tool_dir,
|
||||
source_path=tool_source.source_path,
|
||||
stored_at=tool_source.stored_at,
|
||||
source_metadata=tool_source.metadata or None,
|
||||
)
|
||||
session.add(model)
|
||||
session.flush()
|
||||
@@ -128,38 +114,31 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
"""Retrieve a tool source by hash."""
|
||||
with self._read_session() as session:
|
||||
model = session.execute(
|
||||
select(ToolSourceModel).where(ToolSourceModel.hash == hash)
|
||||
select(ToolSourceRecord).where(ToolSourceRecord.hash == hash)
|
||||
).scalar_one_or_none()
|
||||
if not model:
|
||||
return None
|
||||
return self._model_to_stored(model)
|
||||
|
||||
def _model_to_stored(self, model: ToolSourceModel) -> StoredToolSource:
|
||||
def _model_to_stored(self, model: ToolSourceRecord) -> StoredToolSource:
|
||||
"""Convert database model to StoredToolSource."""
|
||||
source_data = model.source or {}
|
||||
|
||||
stored_at = source_data.get("stored_at")
|
||||
if stored_at and isinstance(stored_at, str):
|
||||
stored_at = datetime.fromisoformat(stored_at)
|
||||
|
||||
assert model.hash is not None
|
||||
return StoredToolSource(
|
||||
hash=model.hash,
|
||||
tool_source_class=source_data.get("tool_source_class", "XmlToolSource"),
|
||||
raw_source=source_data.get("raw", ""),
|
||||
tool_id=source_data.get("tool_id"),
|
||||
tool_version=source_data.get("tool_version"),
|
||||
tool_dir=source_data.get("tool_dir"),
|
||||
source_path=source_data.get("source_path"),
|
||||
stored_at=stored_at,
|
||||
metadata=source_data.get("metadata", {}),
|
||||
tool_source_class=model.source_class or "XmlToolSource",
|
||||
raw_source=model.source or "",
|
||||
tool_id=model.tool_id,
|
||||
tool_version=model.tool_version,
|
||||
tool_dir=model.tool_dir,
|
||||
source_path=model.source_path,
|
||||
stored_at=model.stored_at,
|
||||
metadata=model.source_metadata or {},
|
||||
)
|
||||
|
||||
def exists(self, hash: str) -> bool:
|
||||
"""Check if a tool source exists."""
|
||||
with self._read_session() as session:
|
||||
result = session.execute(
|
||||
select(ToolSourceModel.id).where(ToolSourceModel.hash == hash)
|
||||
select(ToolSourceRecord.id).where(ToolSourceRecord.hash == hash)
|
||||
).scalar_one_or_none()
|
||||
return result is not None
|
||||
|
||||
@@ -167,7 +146,7 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
"""Delete a tool source by hash."""
|
||||
session = self._get_session()
|
||||
|
||||
model = session.execute(select(ToolSourceModel).where(ToolSourceModel.hash == hash)).scalar_one_or_none()
|
||||
model = session.execute(select(ToolSourceRecord).where(ToolSourceRecord.hash == hash)).scalar_one_or_none()
|
||||
|
||||
if not model:
|
||||
return False
|
||||
@@ -182,7 +161,7 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
# yield — otherwise an outer caller could keep the session open
|
||||
# indefinitely while iterating.
|
||||
with self._read_session() as session:
|
||||
result = session.execute(select(ToolSourceModel.hash)).all()
|
||||
result = session.execute(select(ToolSourceRecord.hash)).all()
|
||||
for (hash_value,) in result:
|
||||
if hash_value:
|
||||
yield hash_value
|
||||
@@ -190,35 +169,28 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
def get_by_tool_id(self, tool_id: str, version: str | None = None) -> list[StoredToolSource]:
|
||||
"""Get tool sources by tool ID and optional version."""
|
||||
with self._read_session() as session:
|
||||
# Query all and filter in Python since tool_id is in JSON
|
||||
result = session.execute(select(ToolSourceModel))
|
||||
sources = []
|
||||
for (model,) in result:
|
||||
source_data = model.source or {}
|
||||
if source_data.get("tool_id") == tool_id:
|
||||
if version is None or source_data.get("tool_version") == version:
|
||||
sources.append(self._model_to_stored(model))
|
||||
return sources
|
||||
stmt = select(ToolSourceRecord).where(ToolSourceRecord.tool_id == tool_id)
|
||||
if version is not None:
|
||||
stmt = stmt.where(ToolSourceRecord.tool_version == version)
|
||||
return [self._model_to_stored(model) for model in session.scalars(stmt)]
|
||||
|
||||
def get_by_source_path(self, source_path: str) -> StoredToolSource | None:
|
||||
"""Get the stored source for a given on-disk file path.
|
||||
|
||||
``source_path`` lives inside the JSON ``source`` blob so this scans the
|
||||
table and filters in Python — same shape as ``get_by_tool_id``. The
|
||||
populator writes one entry per file, so there is at most one match.
|
||||
The populator writes one entry per file, so there is at most one match.
|
||||
"""
|
||||
with self._read_session() as session:
|
||||
result = session.execute(select(ToolSourceModel))
|
||||
for (model,) in result:
|
||||
source_data = model.source or {}
|
||||
if source_data.get("source_path") == source_path:
|
||||
return self._model_to_stored(model)
|
||||
return None
|
||||
model = session.execute(
|
||||
select(ToolSourceRecord).where(ToolSourceRecord.source_path == source_path).limit(1)
|
||||
).scalar_one_or_none()
|
||||
if not model:
|
||||
return None
|
||||
return self._model_to_stored(model)
|
||||
|
||||
def count(self) -> int:
|
||||
"""Return the total number of stored tool sources."""
|
||||
with self._read_session() as session:
|
||||
result = session.execute(select(func.count(ToolSourceModel.id)))
|
||||
result = session.execute(select(func.count(ToolSourceRecord.id)))
|
||||
return result.scalar() or 0
|
||||
|
||||
def get_stats(self) -> dict:
|
||||
@@ -290,18 +262,6 @@ class DatabaseToolSourceStore(ToolSourceStore):
|
||||
except Exception as e:
|
||||
log.warning(f"Failed to load index from tool_index table: {e}")
|
||||
|
||||
# Fall back to legacy storage in tool_source table
|
||||
legacy = session.execute(
|
||||
select(ToolSourceModel).where(ToolSourceModel.hash == "__tool_index__")
|
||||
).scalar_one_or_none()
|
||||
|
||||
if legacy:
|
||||
source_data = legacy.source or {}
|
||||
index_data = source_data.get("index")
|
||||
if index_data:
|
||||
self._cached_index = ToolIndex.from_dict(index_data)
|
||||
return self._cached_index
|
||||
|
||||
return None
|
||||
|
||||
def update_index_entry(self, entry: ToolIndexEntry) -> None:
|
||||
|
||||
Reference in New Issue
Block a user