Merge pull request #11914 from Nerdinacan/history-audit

[21.05] Create history_audit table to avoid write deadlocks during history.update_time updates
This commit is contained in:
Marius van den Beek
2021-05-13 16:26:59 +02:00
committed by GitHub
25 changed files with 514 additions and 31 deletions
+14
View File
@@ -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
+1 -6
View File
@@ -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
+11
View File
@@ -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``
~~~~~~~~~~~~~
+11
View File
@@ -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)
+17
View File
@@ -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__':
+10
View File
@@ -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}")
@@ -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'
+1 -1
View File
@@ -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()
+23
View File
@@ -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 = []
+16 -5
View File
@@ -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)),
@@ -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
@@ -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}"
@@ -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()
@@ -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,
)
@@ -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)
@@ -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)
+51
View File
@@ -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)
+5 -2
View File
@@ -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
@@ -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
@@ -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):
@@ -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:
+4 -4
View File
@@ -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))
+40
View File
@@ -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")
+28
View File
@@ -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