Merge pull request #11849 from mvdbeek/restore_update_time_query

Fix collection state updates by setting collection update time
This commit is contained in:
John Chilton
2021-04-19 14:41:39 -04:00
committed by GitHub
6 changed files with 105 additions and 16 deletions
+3 -3
View File
@@ -1321,7 +1321,7 @@ class JobWrapper(HasResourceParameters):
# Pause any dependent jobs (and those jobs' outputs)
for dep_job_assoc in dataset.dependent_jobs:
self.pause(dep_job_assoc.job, "Execution of this dataset's job is paused because its input datasets are in an error state.")
job.set_final_state(job.states.ERROR)
job.set_final_state(job.states.ERROR, supports_skip_locked=self.app.application_stack.supports_skip_locked())
job.command_line = unicodify(self.command_line)
job.info = message
# TODO: Put setting the stdout, stderr, and exit code in one place
@@ -1408,7 +1408,7 @@ class JobWrapper(HasResourceParameters):
job.info = info
job.set_state(state)
self.sa_session.add(job)
job.update_output_states()
job.update_output_states(self.app.application_stack.supports_skip_locked())
if flush:
self.sa_session.flush()
@@ -1776,7 +1776,7 @@ class JobWrapper(HasResourceParameters):
# Finally set the job state. This should only happen *after* all
# dataset creation, and will allow us to eliminate force_history_refresh.
job.set_final_state(final_job_state)
job.set_final_state(final_job_state, supports_skip_locked=self.app.application_stack.supports_skip_locked())
if not job.tasks:
# If job was composed of tasks, don't attempt to recollect statistics
self._collect_metrics(job, job_metrics_directory)
+1 -1
View File
@@ -924,7 +924,7 @@ class JobHandlerStopQueue(Monitors):
if error_msg is not None:
final_state = job.states.ERROR
job.info = error_msg
job.set_final_state(final_state)
job.set_final_state(final_state, supports_skip_locked=self.app.application_stack.supports_skip_locked())
self.sa_session.add(job)
self.sa_session.flush()
+1 -4
View File
@@ -364,12 +364,9 @@ class HistoryContentsManager(containers.ContainerManagerMixin):
# TODO: should be purgable? fix
purged=literal(False),
extension=literal(None),
# these are attached instead to the inner collection joined below
create_time=model.DatasetCollection.create_time,
update_time=model.DatasetCollection.update_time
)
subquery = self._session().query(*columns)
# for the HDCA's we need to join the DatasetCollection since it has update/create times
# for the HDCA's we need to join the DatasetCollection since it has the populated_state
subquery = subquery.join(model.DatasetCollection,
model.DatasetCollection.id == component_class.collection_id)
if history_id:
+38 -4
View File
@@ -37,6 +37,7 @@ from sqlalchemy import (
true,
type_coerce,
types)
from sqlalchemy.exc import OperationalError
from sqlalchemy.ext import hybrid
from sqlalchemy.orm import (
aliased,
@@ -1177,7 +1178,36 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
return rval
def set_final_state(self, final_state):
def update_hdca_update_time_for_job(self, update_time, sa_session, supports_skip_locked):
subq = sa_session.query(HistoryDatasetCollectionAssociation.id) \
.join(ImplicitCollectionJobs) \
.join(ImplicitCollectionJobsJobAssociation) \
.filter(ImplicitCollectionJobsJobAssociation.job_id == self.id)
if supports_skip_locked:
subq = subq.with_for_update(skip_locked=True).subquery()
implicit_statement = HistoryDatasetCollectionAssociation.table.update() \
.where(HistoryDatasetCollectionAssociation.table.c.id.in_(subq)) \
.values(update_time=update_time)
explicit_statement = HistoryDatasetCollectionAssociation.table.update() \
.where(HistoryDatasetCollectionAssociation.table.c.job_id == self.id) \
.values(update_time=update_time)
sa_session.execute(explicit_statement)
if supports_skip_locked:
sa_session.execute(implicit_statement)
else:
conn = sa_session.connection(execution_options={'isolation_level': 'SERIALIZABLE'})
with conn.begin() as trans:
try:
conn.execute(implicit_statement)
trans.commit()
except OperationalError as e:
# If this is a serialization failure on PostgreSQL, then e.orig is a psycopg2 TransactionRollbackError
# and should have attribute `code`. Other engines should just report the message and move on.
if int(getattr(e.orig, 'pgcode', -1)) != 40001:
log.debug(f"Updating implicit collection uptime_time for job {self.id} failed (this is expected for large collections and not a problem): {unicodify(e)}")
trans.rollback()
def set_final_state(self, final_state, supports_skip_locked):
self.set_state(final_state)
# TODO: migrate to where-in subqueries?
statement = '''
@@ -1186,9 +1216,11 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
WHERE job_id = :job_id;
'''
sa_session = object_session(self)
update_time = galaxy.model.orm.now.now()
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked)
params = {
'job_id': self.id,
'update_time': galaxy.model.orm.now.now()
'update_time': update_time
}
sa_session.execute(statement, params)
@@ -1212,7 +1244,7 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
for dataset_assoc in self.output_datasets:
return dataset_assoc.dataset.tool_version
def update_output_states(self):
def update_output_states(self, supports_skip_locked):
# TODO: migrate to where-in subqueries?
statements = ['''
UPDATE dataset
@@ -1256,11 +1288,13 @@ class Job(JobLike, UsesCreateAndUpdateTime, Dictifiable, RepresentById):
);
''']
sa_session = object_session(self)
update_time = galaxy.model.orm.now.now()
self.update_hdca_update_time_for_job(update_time=update_time, sa_session=sa_session, supports_skip_locked=supports_skip_locked)
params = {
'job_id': self.id,
'state': self.state,
'info': self.info,
'update_time': galaxy.model.orm.now.now()
'update_time': update_time
}
for statement in statements:
sa_session.execute(statement, params)
@@ -1,4 +1,6 @@
import json
import time
from datetime import datetime
from requests import delete, put
@@ -563,3 +565,59 @@ class HistoryContentsApiTestCase(ApiTestCase):
assert isinstance(c, dict)
assert 'job_state_summary' in c
assert isinstance(c['job_state_summary'], dict)
def _get_content(self, history_id, update_time):
return self._get(f"/api/histories/{history_id}/contents/near/100/100?update_time-ge={update_time}").json()
def test_history_contents_near_with_update_time(self):
with self.dataset_populator.test_history() as history_id:
first_time = datetime.utcnow().isoformat()
assert len(self._get_content(history_id, update_time=first_time)) == 0
self.dataset_collection_populator.create_list_in_history(history_id=history_id)
assert len(self._get_content(history_id, update_time=first_time)) == 4 # 3 datasets
self.dataset_populator.wait_for_history(history_id)
all_datasets_finished = first_time = datetime.utcnow().isoformat()
assert len(self._get_content(history_id, update_time=all_datasets_finished)) == 0
@skip_without_tool('cat_data_and_sleep')
def test_history_contents_near_with_update_time_implicit_collection(self):
with self.dataset_populator.test_history() as history_id:
hdca_id = self.dataset_collection_populator.create_list_in_history(history_id=history_id).json()['id']
self.dataset_populator.wait_for_history(history_id)
inputs = {
"input1": {'batch': True, 'values': [{"src": "hdca", "id": hdca_id}]},
"sleep_time": 2,
}
response = self.dataset_populator.run_tool(
"cat_data_and_sleep",
inputs,
history_id,
assert_ok=False,
).json()
collection_id = response['implicit_collections'][0]['id']
for _ in range(20):
update_time = datetime.utcnow().isoformat()
time.sleep(1)
update = self._get_content(history_id, update_time=update_time)
if any((c for c in update if c['history_content_type'] == 'dataset_collection' and c['job_state_summary']['ok'] == 3)):
return
raise Exception(f"History content update time query did not include final update for implicit collection {collection_id}")
@skip_without_tool('collection_creates_dynamic_nested')
def test_history_contents_near_with_update_time_explicit_collection(self):
with self.dataset_populator.test_history() as history_id:
inputs = {'foo': 'bar', 'sleep_time': 2}
response = self.dataset_populator.run_tool(
"collection_creates_dynamic_nested",
inputs,
history_id,
assert_ok=False,
).json()
collection_id = response['output_collections'][0]['id']
for _ in range(20):
update_time = datetime.utcnow().isoformat()
time.sleep(1)
update = self._get_content(history_id, update_time=update_time)
if any((c for c in update if c['history_content_type'] == 'dataset_collection' and c['populated_state'] == 'ok')):
return
raise Exception(f"History content update time query did not include populated_state update for dynamic nested collection {collection_id}")
@@ -241,15 +241,15 @@ class HistoryAsContainerTestCase(HistoryAsContainerBaseTestCase):
contents.append(self.add_list_collection_to_history(history, contents[4:6]))
self.log("should allow filtering by update_time")
# in the case of collections we have to change the collection.collection (ugh) to change the update_time
contents[3].collection.populated_state = 'big ball of mud'
# change the update_time by updating the name
contents[3].name = 'big ball of mud'
self.app.model.context.flush()
update_time = contents[3].collection.update_time
update_time = contents[3].update_time
def get_update_time(item):
update_time = getattr(item, 'update_time', None)
if not update_time:
update_time = item.collection.update_time
update_time = item.update_time
return update_time
results = self.contents_manager.contents(history, filters=[parsed_filter("orm", column('update_time') >= update_time)])