diff --git a/config/celery.ini b/config/celery.ini new file mode 100644 index 00000000000..53a358c5843 --- /dev/null +++ b/config/celery.ini @@ -0,0 +1,14 @@ +[env] +PYTHONPATH = lib + +[watcher:celery] +cmd = celery +args = --app galaxy.celery worker --concurrency 2 -l debug +copy_env = True +numprocesses = 1 + +[watcher:celery-beat] +cmd = celery +args = --app galaxy.celery beat -l debug +copy_env = True +numprocesses = 1 diff --git a/config/dev.ini b/config/dev.ini index d4e5494fe7b..56bd44ee53f 100644 --- a/config/dev.ini +++ b/config/dev.ini @@ -1,5 +1,6 @@ [circus] debug = True +include = celery.ini [env] PYTHONPATH=lib @@ -25,9 +26,3 @@ stop_children = True [socket:web] host = 0.0.0.0 port = 8080 - -[watcher:celery] -cmd = celery -args = --app galaxy.celery worker -l debug -copy_env = True -numprocesses = 1 diff --git a/doc/source/admin/galaxy_options.rst b/doc/source/admin/galaxy_options.rst index f14f01aed1a..0b03370ee25 100644 --- a/doc/source/admin/galaxy_options.rst +++ b/doc/source/admin/galaxy_options.rst @@ -246,6 +246,17 @@ :Type: float +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ +``history_audit_table_prune_interval`` +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +:Description: + Time (in seconds) between attempts to remove old rows from the + history_audit database table. Set to 0 to disable pruning. +:Default: ``3600`` +:Type: int + + ~~~~~~~~~~~~~ ``file_path`` ~~~~~~~~~~~~~ diff --git a/lib/galaxy/app.py b/lib/galaxy/app.py index fa6b7e6b62a..dffa996a080 100644 --- a/lib/galaxy/app.py +++ b/lib/galaxy/app.py @@ -61,6 +61,7 @@ from galaxy.util import ( heartbeat, StructuredExecutionTimer, ) +from galaxy.util.task import IntervalTask from galaxy.visualization.data_providers.registry import DataProviderRegistry from galaxy.visualization.genomes import Genomes from galaxy.visualization.plugins.registry import VisualizationsRegistry @@ -304,6 +305,16 @@ class UniverseApplication(StructuredApp, GalaxyManagerApplication): self.authnz_manager = managers.AuthnzManager(self, self.config.oidc_config_file, self.config.oidc_backends_config_file) + + if not self.config.enable_celery_tasks and self.config.history_audit_table_prune_interval > 0: + self.prune_history_audit_task = IntervalTask( + func=lambda: galaxy.model.HistoryAudit.prune(self.model.session), + name="HistoryAuditTablePruneTask", + interval=self.config.history_audit_table_prune_interval, + immediate_start=False, + time_execution=True) + self.application_stack.register_postfork_function(self.prune_history_audit_task.start) + self.haltables.append(("HistoryAuditTablePruneTask", self.prune_history_audit_task.shutdown)) # Start the job manager self.application_stack.register_postfork_function(self.job_manager.start) self.proxy_manager = ProxyManager(self.config) diff --git a/lib/galaxy/celery/__init__.py b/lib/galaxy/celery/__init__.py index d8f916e3738..4d7c573ecdf 100644 --- a/lib/galaxy/celery/__init__.py +++ b/lib/galaxy/celery/__init__.py @@ -52,8 +52,25 @@ def get_broker(): return config.amqp_internal_connection +def get_history_audit_table_prune_interval(): + config = get_config() + if config: + return config.history_audit_table_prune_interval + else: + return 3600 + + broker = get_broker() celery_app = Celery('galaxy', broker=broker, include=['galaxy.celery.tasks']) +prune_interval = get_history_audit_table_prune_interval() +if prune_interval > 0: + celery_app.conf.beat_schedule = { + 'prune-history-audit-table': { + 'task': 'galaxy.celery.tasks.prune_history_audit_table', + 'schedule': prune_interval, + }, + } +celery_app.conf.timezone = 'UTC' if __name__ == '__main__': diff --git a/lib/galaxy/celery/tasks.py b/lib/galaxy/celery/tasks.py index 878c080f64b..a5238e8b7c1 100644 --- a/lib/galaxy/celery/tasks.py +++ b/lib/galaxy/celery/tasks.py @@ -9,6 +9,7 @@ from galaxy.celery import celery_app from galaxy.jobs.manager import JobManager from galaxy.managers.hdas import HDAManager from galaxy.managers.lddas import LDDAManager +from galaxy.util import ExecutionTimer from galaxy.util.custom_logging import get_logger from . import get_galaxy_app @@ -73,3 +74,12 @@ def export_history( job.state = model.Job.states.NEW sa_session.flush() job_manager.enqueue(job) + + +@celery_app.task +@galaxy_task +def prune_history_audit_table(sa_session: scoped_session): + """Prune ever growing history_audit table.""" + timer = ExecutionTimer() + model.HistoryAudit.prune(sa_session) + log.debug(f"Successfully pruned history_audit table {timer}") diff --git a/lib/galaxy/config/sample/galaxy.yml.sample b/lib/galaxy/config/sample/galaxy.yml.sample index f9026c7042e..6349fda6778 100644 --- a/lib/galaxy/config/sample/galaxy.yml.sample +++ b/lib/galaxy/config/sample/galaxy.yml.sample @@ -222,6 +222,10 @@ galaxy: # seconds). #database_wait_sleep: 1.0 + # Time (in seconds) between attempts to remove old rows from the + # history_audit database table. Set to 0 to disable pruning. + #history_audit_table_prune_interval: 3600 + # Where dataset files are stored. It must be accessible at the same # path on any cluster nodes that will run Galaxy jobs, unless using # Pulsar. The default value has been changed from 'files' to 'objects' diff --git a/lib/galaxy/managers/histories.py b/lib/galaxy/managers/histories.py index 534d6f123dc..026f2155579 100644 --- a/lib/galaxy/managers/histories.py +++ b/lib/galaxy/managers/histories.py @@ -85,7 +85,7 @@ class HistoryManager(sharable.SharableModelManager, deletable.PurgableManagerMix """ if self.user_manager.is_anonymous(user): return None if (not current_history or current_history.deleted) else current_history - desc_update_time = desc(self.model_class.table.c.update_time) + desc_update_time = desc(self.model_class.update_time) filters = self._munge_filters(filters, self.model_class.user_id == user.id) # TODO: normalize this return value return self.query(filters=filters, order_by=desc_update_time, limit=1, **kwargs).first() diff --git a/lib/galaxy/model/__init__.py b/lib/galaxy/model/__init__.py index 873edab4fb4..af5c156ecf2 100644 --- a/lib/galaxy/model/__init__.py +++ b/lib/galaxy/model/__init__.py @@ -35,6 +35,7 @@ from sqlalchemy import ( select, text, true, + tuple_, type_coerce, types) from sqlalchemy.exc import OperationalError @@ -1845,6 +1846,27 @@ def is_hda(d): return isinstance(d, HistoryDatasetAssociation) +class HistoryAudit(RepresentById): + def __init__(self, history, update_time): + self.history = history + self.update_time = update_time + + @classmethod + def prune(cls, sa_session): + history_audit_table = cls.table + latest_subq = sa_session.query( + history_audit_table.c.history_id, + func.max(history_audit_table.c.update_time).label('max_update_time')).group_by(history_audit_table.c.history_id).subquery() + not_latest_query = sa_session.query( + history_audit_table.c.history_id, history_audit_table.c.update_time + ).select_from(latest_subq).join( + history_audit_table, and_( + history_audit_table.c.update_time < latest_subq.columns.max_update_time, + history_audit_table.c.history_id == latest_subq.columns.history_id)) + d = history_audit_table.delete() + sa_session.execute(d.where(tuple_(history_audit_table.c.history_id, history_audit_table.c.update_time).in_(not_latest_query))) + + class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): dict_collection_visible_keys = ['id', 'name', 'published', 'deleted'] @@ -1860,6 +1882,7 @@ class History(HasTags, Dictifiable, UsesAnnotations, HasName, RepresentById): self.importing = False self.genome_build = None self.published = False + self.update_time = None # Relationships self.user = user self.datasets = [] diff --git a/lib/galaxy/model/mapping.py b/lib/galaxy/model/mapping.py index e8fc136bf89..bbffbd36d5d 100644 --- a/lib/galaxy/model/mapping.py +++ b/lib/galaxy/model/mapping.py @@ -22,6 +22,7 @@ from sqlalchemy import ( MetaData, not_, Numeric, + PrimaryKeyConstraint, select, String, Table, TEXT, @@ -46,10 +47,10 @@ from galaxy.model.custom_types import ( TrimmedString, UUIDType, ) +from galaxy.model.migrate.triggers.update_audit_table import install as install_timestamp_triggers from galaxy.model.orm.engine_factory import build_engine from galaxy.model.orm.now import now from galaxy.model.security import GalaxyRBACAgent -from galaxy.model.triggers import install_timestamp_triggers from galaxy.model.view import HistoryDatasetCollectionJobStateSummary from galaxy.model.view.utils import install_views @@ -202,7 +203,7 @@ model.History.table = Table( "history", metadata, Column("id", Integer, primary_key=True), Column("create_time", DateTime, default=now), - Column("update_time", DateTime, index=True, default=now, onupdate=now), + Column("update_time", DateTime, key="_update_time", index=True, default=now, onupdate=now), Column("user_id", Integer, ForeignKey("galaxy_user.id"), index=True), Column("name", TrimmedString(255)), Column("hid_counter", Integer, default=1), @@ -216,6 +217,13 @@ model.History.table = Table( Index('ix_history_slug', 'slug', mysql_length=200), ) +model.HistoryAudit.table = Table( + "history_audit", metadata, + Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), + Column("update_time", DateTime, default=now, primary_key=True, nullable=False), + PrimaryKeyConstraint(sqlite_on_conflict='IGNORE') +) + model.HistoryUserShareAssociation.table = Table( "history_user_share_association", metadata, Column("id", Integer, primary_key=True), @@ -1889,7 +1897,10 @@ mapper(model.History, model.History.table, properties=dict( users_shared_with_count=column_property( select([func.count(model.HistoryUserShareAssociation.table.c.id)]).where(model.History.table.c.id == model.HistoryUserShareAssociation.table.c.history_id), deferred=True - ) + ), + update_time=column_property( + select([func.max(model.HistoryAudit.table.c.update_time)]).where(model.HistoryAudit.table.c.history_id == model.History.table.c.id), + ), )) # Set up proxy so that @@ -1905,13 +1916,13 @@ mapper(model.HistoryUserShareAssociation, model.HistoryUserShareAssociation.tabl mapper(model.User, model.User.table, properties=dict( histories=relation(model.History, backref="user", - order_by=desc(model.History.table.c.update_time)), + order_by=desc(model.History.update_time)), active_histories=relation(model.History, primaryjoin=( (model.History.table.c.user_id == model.User.table.c.id) & (not_(model.History.table.c.deleted)) ), - order_by=desc(model.History.table.c.update_time)), + order_by=desc(model.History.update_time)), galaxy_sessions=relation(model.GalaxySession, order_by=desc(model.GalaxySession.table.c.update_time)), diff --git a/lib/galaxy/model/migrate/triggers/__init__.py b/lib/galaxy/model/migrate/triggers/__init__.py new file mode 100644 index 00000000000..e69de29bb2d diff --git a/lib/galaxy/model/triggers.py b/lib/galaxy/model/migrate/triggers/history_update_time_field.py similarity index 97% rename from lib/galaxy/model/triggers.py rename to lib/galaxy/model/migrate/triggers/history_update_time_field.py index 8464c9822d6..3cbfa83eb9f 100644 --- a/lib/galaxy/model/triggers.py +++ b/lib/galaxy/model/migrate/triggers/history_update_time_field.py @@ -2,7 +2,7 @@ Database trigger installation and removal """ -from sqlalchemy import DDL +from galaxy.model.migrate.versions.util import execute_statements def install_timestamp_triggers(engine): @@ -21,12 +21,6 @@ def drop_timestamp_triggers(engine): execute_statements(engine, statements) -def execute_statements(engine, statements): - for sql in statements: - cmd = DDL(sql) - cmd.execute(bind=engine) - - def get_timestamp_install_sql(variant): """ Generate a list of SQL statements for installation of timestamp triggers diff --git a/lib/galaxy/model/migrate/triggers/update_audit_table.py b/lib/galaxy/model/migrate/triggers/update_audit_table.py new file mode 100644 index 00000000000..d3340f8b727 --- /dev/null +++ b/lib/galaxy/model/migrate/triggers/update_audit_table.py @@ -0,0 +1,172 @@ +from galaxy.model.migrate.versions.util import execute_statements + + +# function name prefix +fn_prefix = "fn_audit_history_by" + +# map between source table and associated incoming id field +trigger_config = { + 'history_dataset_association': "history_id", + 'history_dataset_collection_association': "history_id", + 'history': "id", +} + + +def install(engine): + """Install history audit table triggers""" + sql = _postgres_install(engine) if 'postgres' in engine.name else _sqlite_install() + execute_statements(engine, sql) + + +def remove(engine): + """Uninstall history audit table triggers""" + sql = _postgres_remove() if 'postgres' in engine.name else _sqlite_remove() + execute_statements(engine, sql) + + +# Postgres trigger installation + + +def _postgres_remove(): + """postgres trigger removal sql""" + + sql = [] + sql.append(f"DROP FUNCTION IF EXISTS {fn_prefix}_history_id() CASCADE;") + sql.append(f"DROP FUNCTION IF EXISTS {fn_prefix}_id() CASCADE;") + + return sql + + +def _postgres_install(engine): + """postgres trigger installation sql""" + + sql = [] + + # postgres trigger function template + # need to make separate functions purely because the incoming history_id field name will be + # different for different source tables. There may be a fancier way to dynamically choose + # between incoming fields, but having 2 triggers fns seems straightforward + + def statement_trigger_fn(id_field): + fn = f"{fn_prefix}_{id_field}" + + return f""" + CREATE OR REPLACE FUNCTION {fn}() + RETURNS TRIGGER + LANGUAGE 'plpgsql' + AS $BODY$ + BEGIN + INSERT INTO history_audit (history_id, update_time) + SELECT DISTINCT {id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC' + FROM new_table + WHERE {id_field} IS NOT NULL + ON CONFLICT DO NOTHING; + RETURN NULL; + END; + $BODY$ + """ + + def row_trigger_fn(id_field): + fn = f"{fn_prefix}_{id_field}" + + return f""" + CREATE OR REPLACE FUNCTION {fn}() + RETURNS TRIGGER + LANGUAGE 'plpgsql' + AS $BODY$ + BEGIN + INSERT INTO history_audit (history_id, update_time) + VALUES (NEW.{id_field}, CURRENT_TIMESTAMP AT TIME ZONE 'UTC') + ON CONFLICT DO NOTHING; + RETURN NULL; + END; + $BODY$ + """ + + def statement_trigger_def(source_table, id_field, operation, when="AFTER"): + fn = f"{fn_prefix}_{id_field}" + + # Postgres supports many triggers per operation/table so the label can + # be indicative of what's happening + label = f"history_audit_by_{id_field}" + trigger_name = get_trigger_name(label, operation, when, statement=True) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} ON {source_table} + REFERENCING NEW TABLE AS new_table + FOR EACH STATEMENT EXECUTE FUNCTION {fn}(); + """ + + def row_trigger_def(source_table, id_field, operation, when="AFTER"): + fn = f"{fn_prefix}_{id_field}" + + label = f"history_audit_by_{id_field}" + trigger_name = get_trigger_name(label, operation, when, statement=True) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} ON {source_table} + FOR EACH ROW EXECUTE FUNCTION {fn}(); + """ + + # pick row or statement triggers depending on postgres version + version = engine.dialect.server_version_info[0] + trigger_fn = statement_trigger_fn if version > 10 else row_trigger_fn + trigger_def = statement_trigger_def if version > 10 else row_trigger_def + + for id_field in ["history_id", "id"]: + sql.append(trigger_fn(id_field)) + + for source_table, id_field in trigger_config.items(): + for operation in ["UPDATE", "INSERT"]: + sql.append(trigger_def(source_table, id_field, operation)) + + return sql + + +def _sqlite_remove(): + sql = [] + + for source_table in trigger_config: + for operation in ["UPDATE", "INSERT"]: + trigger_name = get_trigger_name(source_table, operation, "AFTER") + sql.append(f"DROP TRIGGER IF EXISTS {trigger_name};") + + return sql + + +def _sqlite_install(): + # delete old stuff first + sql = _sqlite_remove() + + def trigger_def(source_table, id_field, operation, when="AFTER"): + + # only one trigger per operation/table in simple databases, so + # trigger name is less descriptive + trigger_name = get_trigger_name(source_table, operation, when) + + return f""" + CREATE TRIGGER {trigger_name} + {when} {operation} + ON {source_table} + FOR EACH ROW + BEGIN + INSERT INTO history_audit (history_id, update_time) + SELECT NEW.{id_field}, strftime('%%Y-%%m-%%d %%H:%%M:%%f', 'now') + WHERE NEW.{id_field} IS NOT NULL; + END; + """ + + for source_table, id_field in trigger_config.items(): + for operation in ["UPDATE", "INSERT"]: + sql.append(trigger_def(source_table, id_field, operation)) + + return sql + + +def get_trigger_name(label, operation, when, statement=False): + op_initial = operation.lower()[0] + when_initial = when.lower()[0] + rs = "s" if statement else "r" + return f"trigger_{label}_{when_initial}{op_initial}{rs}" diff --git a/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py b/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py index dc7b0f9647c..607dee1150d 100644 --- a/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py +++ b/lib/galaxy/model/migrate/versions/0165_add_content_update_time.py @@ -7,9 +7,9 @@ import logging from sqlalchemy import Column, DateTime, MetaData, Table +from galaxy.model.migrate.triggers.history_update_time_field import drop_timestamp_triggers, install_timestamp_triggers from galaxy.model.migrate.versions.util import add_column, drop_column from galaxy.model.orm.now import now -from galaxy.model.triggers import drop_timestamp_triggers, install_timestamp_triggers log = logging.getLogger(__name__) metadata = MetaData() diff --git a/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py b/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py index 031ce925ef3..bd1ad2931c8 100644 --- a/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py +++ b/lib/galaxy/model/migrate/versions/0174_readd_update_time_triggers.py @@ -6,7 +6,7 @@ import logging from sqlalchemy import MetaData -from galaxy.model.triggers import ( +from galaxy.model.migrate.triggers.history_update_time_field import ( drop_timestamp_triggers, install_timestamp_triggers, ) diff --git a/lib/galaxy/model/migrate/versions/0175_history_audit.py b/lib/galaxy/model/migrate/versions/0175_history_audit.py new file mode 100644 index 00000000000..521bb762396 --- /dev/null +++ b/lib/galaxy/model/migrate/versions/0175_history_audit.py @@ -0,0 +1,83 @@ +""" +Add history audit table and associated triggers +""" + +import datetime +import logging + +from sqlalchemy import Column, DateTime, ForeignKey, Integer, MetaData, PrimaryKeyConstraint, Table + +from galaxy.model.migrate.triggers import ( + history_update_time_field as old_triggers, # rollback to old ones + update_audit_table as new_triggers, # install me +) +from galaxy.model.migrate.versions.util import ( + create_table, + drop_table +) + +log = logging.getLogger(__name__) +now = datetime.datetime.utcnow +metadata = MetaData() + +AuditTable = Table( + "history_audit", + metadata, + Column("history_id", Integer, ForeignKey("history.id"), primary_key=True, nullable=False), + Column("update_time", DateTime, default=now, primary_key=True, nullable=False), + PrimaryKeyConstraint(sqlite_on_conflict='IGNORE') +) + + +def upgrade(migrate_engine): + print(__doc__) + metadata.bind = migrate_engine + metadata.reflect() + + # create table + index + AuditTable.drop(migrate_engine, checkfirst=True) + create_table(AuditTable) + + # populate with update_time from every history + copy_update_times = """ + INSERT INTO history_audit (history_id, update_time) + SELECT id, update_time FROM history + """ + migrate_engine.execute(copy_update_times) + + # drop existing timestamp triggers + old_triggers.drop_timestamp_triggers(migrate_engine) + + # install new timestamp triggers + new_triggers.install(migrate_engine) + + +def downgrade(migrate_engine): + print(__doc__) + metadata.bind = migrate_engine + metadata.reflect() + + # drop existing timestamp triggers + new_triggers.remove(migrate_engine) + + try: + # update history.update_time with vals from audit table + put_em_back = """ + UPDATE history h + SET update_time = a.max_update_time + FROM ( + SELECT history_id, max(update_time) as max_update_time + FROM history_audit + GROUP BY history_id + ) a + WHERE h.id = a.history_id + """ + migrate_engine.execute(put_em_back) + except Exception: + print("Unable to put update_times back") + + # drop audit table + drop_table(AuditTable) + + # install old timestamp triggers + old_triggers.install_timestamp_triggers(migrate_engine) diff --git a/lib/galaxy/model/migrate/versions/util.py b/lib/galaxy/model/migrate/versions/util.py index a9265f9a5e1..e961e07b4c2 100644 --- a/lib/galaxy/model/migrate/versions/util.py +++ b/lib/galaxy/model/migrate/versions/util.py @@ -3,6 +3,7 @@ import logging from sqlalchemy import ( BLOB, + DDL, Index, Table, Text @@ -193,3 +194,10 @@ def drop_index(index, table, column_name=None, metadata=None): index.drop() except Exception: log.exception("Dropping index '%s' from table '%s' failed", index, table) + + +def execute_statements(engine, raw_sql): + statements = raw_sql if isinstance(raw_sql, list) else [raw_sql] + for sql in statements: + cmd = DDL(sql) + cmd.execute(bind=engine) diff --git a/lib/galaxy/util/task.py b/lib/galaxy/util/task.py new file mode 100644 index 00000000000..9da0e347d08 --- /dev/null +++ b/lib/galaxy/util/task.py @@ -0,0 +1,51 @@ +import logging +from threading import ( + Event, + Thread, +) + +from galaxy.util import ExecutionTimer + +log = logging.getLogger(__name__) + + +class IntervalTask: + + def __init__(self, func, name="Periodic task", interval=3600, immediate_start=False, time_execution=False): + """ + Run an arbitrary function `func` every `interval` seconds. + + Set `immediate_start` to True to run `func` when task is started. + """ + self.func = func + self.name = name + self.interval = interval + self.time_execution = time_execution + self.immediate_start = immediate_start + self.event = Event() + self.thread = Thread(target=self.run, name=self.name, daemon=True) + self.running = False + + def start(self): + self.running = True + self.thread.start() + + def _exec(self): + if self.time_execution: + timer = ExecutionTimer() + self.func() + if self.time_execution: + log.debug(f"Executed periodic task {self.name} {timer}") + + def run(self): + if self.immediate_start: + self._exec() + while not self.event.isSet(): + self.event.wait(self.interval) + if self.running: + self._exec() + + def shutdown(self): + self.running = False + self.event.set() + self.thread.join(5) diff --git a/lib/galaxy/web/framework/helpers/grids.py b/lib/galaxy/web/framework/helpers/grids.py index 4568987d3ac..f741e931bf4 100644 --- a/lib/galaxy/web/framework/helpers/grids.py +++ b/lib/galaxy/web/framework/helpers/grids.py @@ -76,10 +76,13 @@ class GridColumn: """Sort query using this column.""" if column_name is None: column_name = self.key + column = self.model_class.table.c.get(column_name) + if column is None: + column = getattr(self.model_class, column_name) if ascending: - query = query.order_by(self.model_class.table.c.get(column_name).asc()) + query = query.order_by(column.asc()) else: - query = query.order_by(self.model_class.table.c.get(column_name).desc()) + query = query.order_by(column.desc()) return query diff --git a/lib/galaxy/webapps/galaxy/config_schema.yml b/lib/galaxy/webapps/galaxy/config_schema.yml index 83ea2df6742..43081fdd4b7 100644 --- a/lib/galaxy/webapps/galaxy/config_schema.yml +++ b/lib/galaxy/webapps/galaxy/config_schema.yml @@ -197,6 +197,14 @@ mapping: desc: | Time to sleep between attempts if database_wait is enabled (in seconds). + history_audit_table_prune_interval: + type: int + default: 3600 + required: false + desc: | + Time (in seconds) between attempts to remove old rows from the history_audit database table. + Set to 0 to disable pruning. + file_path: type: str default: objects diff --git a/lib/galaxy/webapps/galaxy/controllers/history.py b/lib/galaxy/webapps/galaxy/controllers/history.py index 96dba4ba74f..77a2b183631 100644 --- a/lib/galaxy/webapps/galaxy/controllers/history.py +++ b/lib/galaxy/webapps/galaxy/controllers/history.py @@ -86,9 +86,9 @@ class HistoryListGrid(grids.Grid): def sort(self, trans, query, ascending, column_name=None): if ascending: - query = query.order_by(self.model_class.table.c.purged.asc(), self.model_class.table.c.update_time.desc()) + query = query.order_by(self.model_class.table.c.purged.asc(), self.model_class.update_time.desc()) else: - query = query.order_by(self.model_class.table.c.purged.desc(), self.model_class.table.c.update_time.desc()) + query = query.order_by(self.model_class.table.c.purged.desc(), self.model_class.update_time.desc()) return query def build_initial_query(self, trans, **kwargs): diff --git a/lib/galaxy/webapps/reports/controllers/system.py b/lib/galaxy/webapps/reports/controllers/system.py index b31a1123d3b..1d56d955075 100644 --- a/lib/galaxy/webapps/reports/controllers/system.py +++ b/lib/galaxy/webapps/reports/controllers/system.py @@ -60,7 +60,7 @@ class System(BaseUIController): for history in trans.sa_session.query(model.History) \ .filter(and_(model.History.table.c.user_id == null(), model.History.table.c.deleted == true(), - model.History.table.c.update_time < cutoff_time)): + model.History.update_time < cutoff_time)): for dataset in history.datasets: if not dataset.deleted: dataset_count += 1 @@ -86,7 +86,7 @@ class System(BaseUIController): histories = trans.sa_session.query(model.History) \ .filter(and_(model.History.table.c.deleted == true(), model.History.table.c.purged == false(), - model.History.table.c.update_time < cutoff_time)) \ + model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) for history in histories: diff --git a/scripts/cleanup_datasets/cleanup_datasets.py b/scripts/cleanup_datasets/cleanup_datasets.py index 9934b5c297b..6f3311d0c1b 100755 --- a/scripts/cleanup_datasets/cleanup_datasets.py +++ b/scripts/cleanup_datasets/cleanup_datasets.py @@ -140,12 +140,12 @@ def delete_userless_histories(app, cutoff_time, info_only=False, force_retry=Fal if force_retry: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.user_id == null(), - app.model.History.table.c.update_time < cutoff_time)) + app.model.History.update_time < cutoff_time)) else: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.user_id == null(), app.model.History.table.c.deleted == false(), - app.model.History.table.c.update_time < cutoff_time)) + app.model.History.update_time < cutoff_time)) for history in histories: if not info_only: log.info("Deleting history id %d", history.id) @@ -170,13 +170,13 @@ def purge_histories(app, cutoff_time, remove_from_disk, info_only=False, force_r if force_retry: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.deleted == true(), - app.model.History.table.c.update_time < cutoff_time)) \ + app.model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) else: histories = app.sa_session.query(app.model.History) \ .filter(and_(app.model.History.table.c.deleted == true(), app.model.History.table.c.purged == false(), - app.model.History.table.c.update_time < cutoff_time)) \ + app.model.History.update_time < cutoff_time)) \ .options(eagerload('datasets')) for history in histories: log.info("### Processing history id %d (%s)", history.id, unicodify(history.name)) diff --git a/test/unit/data/test_galaxy_mapping.py b/test/unit/data/test_galaxy_mapping.py index 94064615a2d..d68a984d45b 100644 --- a/test/unit/data/test_galaxy_mapping.py +++ b/test/unit/data/test_galaxy_mapping.py @@ -472,6 +472,46 @@ class MappingTests(BaseModelTestCase): assert contents_iter_names(ids=[d1.id, d3.id]) == ["1", "3"] + def test_history_audit(self): + model = self.model + u = model.User(email="contents@foo.bar.baz", password="password") + h1 = model.History(name="HistoryAuditHistory", user=u) + h2 = model.History(name="HistoryAuditHistory", user=u) + + def get_audit_table_entries(history): + return self.session().query(model.HistoryAudit.table).filter( + model.HistoryAudit.table.c.history_id == history.id).all() + + def get_latest_entry(entries): + # key ensures result is correct if new columns are added + return max(entries, key=lambda x: x.update_time) + + self.persist(u, h1, h2, expunge=False) + assert len(get_audit_table_entries(h1)) == 1 + assert len(get_audit_table_entries(h2)) == 1 + + self.new_hda(h1, name="1") + self.new_hda(h2, name="2") + self.session().flush() + # db_next_hid modifies history, plus trigger on HDA means 2 additional audit rows per history + + h1_audits = get_audit_table_entries(h1) + h2_audits = get_audit_table_entries(h2) + assert len(h1_audits) == 3 + assert len(h2_audits) == 3 + + h1_latest = get_latest_entry(h1_audits) + h2_latest = get_latest_entry(h2_audits) + + model.HistoryAudit.prune(self.session()) + + h1_audits = get_audit_table_entries(h1) + h2_audits = get_audit_table_entries(h2) + assert len(h1_audits) == 1 + assert len(h2_audits) == 1 + assert h1_audits[0] == h1_latest + assert h2_audits[0] == h2_latest + def _non_empty_flush(self): model = self.model lf = model.LibraryFolder(name="RootFolder") diff --git a/test/unit/util/test_task.py b/test/unit/util/test_task.py new file mode 100644 index 00000000000..a1b7b317634 --- /dev/null +++ b/test/unit/util/test_task.py @@ -0,0 +1,28 @@ +import time + +from galaxy.util.task import IntervalTask + + +def test_interval_task_immediate_start(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=True) + task.start() + task.shutdown() + assert len(results) == 1 + + +def test_interval_task_delayed_start(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=False) + task.start() + task.shutdown() + assert len(results) == 0 + + +def test_interval_task_delayed_start_run_once(): + results = [] + task = IntervalTask(lambda: results.append(1), name="test_task", interval=0.2, immediate_start=False) + task.start() + time.sleep(0.25) + task.shutdown() + assert len(results) == 1