Deduplicate pgcleanup.py

This commit is contained in:
Nate Coraor
2018-09-18 18:23:34 -04:00
parent 5ace7d4015
commit eae35fab33
+124 -137
View File
@@ -15,6 +15,7 @@ import shutil
import sys
import psycopg2
from six import string_types
from sqlalchemy.engine.url import make_url
galaxy_root = os.path.abspath(os.path.join(os.path.dirname(__file__), os.pardir, os.pardir))
@@ -161,18 +162,18 @@ class Cleanup(object):
"""
if message is None:
message = inspect.stack()[1][3]
message = inspect.stack()[2][3]
sql = """
INSERT INTO cleanup_event
(create_time, message)
VALUES (NOW(), %s)
VALUES (NOW(), %(message)s)
RETURNING id;
"""
log.debug("SQL is: %s" % sql % ("'" + message + "'"))
args = {'message': message}
args = (message,)
log.debug("SQL is: %s", sql % args_for_print(args))
cur = self.conn.cursor()
@@ -197,11 +198,40 @@ class Cleanup(object):
return cur.fetchone()[0]
def _update(self, sql, args):
@property
def _update_time_sql(self):
# only update time if not doing force retry (otherwise a lot of things would have their update times reset that were actually purged a long time ago)
if self.args.update_time and not self.args.force_retry:
return ", update_time = NOW()"
return ""
@property
def _force_retry_sql(self):
if not self.args.force_retry:
return " AND NOT purged"""
return ""
def _format_sql(self, sql):
return sql.format(
update_time_sql=self._update_time_sql,
force_retry_sql=self._force_retry_sql,
)
def _add_default_args(self, args):
if isinstance(args, dict):
if 'days' not in args:
args['days'] = self.args.days
def _update(self, sql, args, add_event=True, event_message=None):
sql = self._format_sql(sql)
if add_event and isinstance(args, dict) and 'event_id' not in args:
args['event_id'] = self._create_event(message=event_message)
self._add_default_args(args)
if args is not None:
log.debug('SQL is: %s' % sql % args)
log.debug('SQL is: %s', sql % args_for_print(args))
else:
log.debug('SQL is: %s' % sql)
log.debug('SQL is: %s', sql)
cur = self.conn.cursor()
@@ -260,19 +290,19 @@ class Cleanup(object):
FROM history_dataset_association hda
JOIN history h ON h.id = hda.history_id
JOIN dataset d ON hda.dataset_id = d.id
WHERE h.user_id = %s
WHERE h.user_id = %(user_id)s
AND h.purged = false
AND hda.purged = false
AND d.purged = false
AND d.id NOT IN (SELECT dataset_id
FROM library_dataset_dataset_association)
GROUP BY d.id) sizes)
WHERE id = %s
WHERE id = %(user_id)s
RETURNING disk_usage;
"""
args = (user_id, user_id)
cur = self._update(sql, args)
args = {'user_id': user_id}
cur = self._update(sql, args, add_event=False)
self._flush()
for tup in cur:
@@ -288,14 +318,21 @@ class Cleanup(object):
handle.write('\n')
handle.close()
#
# actions
#
# the _update method automatically:
# - calls the sql argument's (string) format() method to replace '{update_time_sql}' and '{force_retry_sql}'
# - creates an event and adds its id to args with key `event_id` (if args is a dict)
# - adds `days` to args (if args is a dict)
def update_hda_purged_flag(self):
"""
The old cleanup script does not mark HistoryDatasetAssociations as purged when deleted Histories are purged. This method can be used to rectify that situation.
"""
log.info('Marking purged all HistoryDatasetAssociations associated with purged Datasets')
event_id = self._create_event()
# update_time is intentionally left unmodified.
sql = """
WITH purged_hda_ids
@@ -309,15 +346,14 @@ class Cleanup(object):
hda_events
AS (INSERT INTO cleanup_event_hda_association
(create_time, cleanup_event_id, hda_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM purged_hda_ids)
SELECT id
FROM purged_hda_ids
ORDER BY id;
"""
args = (event_id,)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -331,34 +367,25 @@ class Cleanup(object):
"""
log.info('Marking deleted all userless Histories older than %i days' % self.args.days)
event_id = self._create_event()
sql = """
WITH deleted_history_ids
AS ( UPDATE history
SET deleted = true%s
SET deleted = true{update_time_sql}
WHERE user_id is null
AND NOT deleted
AND update_time < (NOW() - interval '%s days')
AND update_time < (NOW() - interval '%(days)s days')
RETURNING id),
history_events
AS (INSERT INTO cleanup_event_history_association
(create_time, cleanup_event_id, history_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_history_ids)
SELECT id
FROM deleted_history_ids
ORDER BY id;
ORDER BY id
"""
update_time_sql = ''
if self.args.update_time:
update_time_sql = """,
update_time = NOW()"""
sql = sql % (update_time_sql, '%s', '%s')
args = (self.args.days, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -375,19 +402,17 @@ class Cleanup(object):
"""
log.info('Marking purged all deleted HistoryDatasetAssociations older than %i days' % self.args.days)
event_id = self._create_event()
sql = """
WITH purged_hda_ids
AS ( UPDATE history_dataset_association
SET purged = true%s
WHERE deleted%s
AND update_time < (NOW() - interval '%s days')
SET purged = true{update_time_sql}
WHERE deleted{force_retry_sql}
AND update_time < (NOW() - interval '%(days)s days')
RETURNING id,
history_id),
deleted_metadata_file_ids
AS ( UPDATE metadata_file
SET deleted = true%s
SET deleted = true{update_time_sql}
FROM purged_hda_ids
WHERE purged_hda_ids.id = metadata_file.hda_id
RETURNING metadata_file.hda_id AS hda_id,
@@ -395,7 +420,7 @@ class Cleanup(object):
metadata_file.object_store_id AS object_store_id),
deleted_icda_ids
AS ( UPDATE implicitly_converted_dataset_association
SET deleted = true%s
SET deleted = true{update_time_sql}
FROM purged_hda_ids
WHERE purged_hda_ids.id = implicitly_converted_dataset_association.hda_parent_id
RETURNING implicitly_converted_dataset_association.hda_id AS hda_id,
@@ -403,28 +428,28 @@ class Cleanup(object):
implicitly_converted_dataset_association.id AS id),
deleted_icda_purged_child_hda_ids
AS ( UPDATE history_dataset_association
SET purged = true%s
SET purged = true{update_time_sql}
FROM deleted_icda_ids
WHERE deleted_icda_ids.hda_id = history_dataset_association.id),
hda_events
AS (INSERT INTO cleanup_event_hda_association
(create_time, cleanup_event_id, hda_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM purged_hda_ids),
metadata_file_events
AS (INSERT INTO cleanup_event_metadata_file_association
(create_time, cleanup_event_id, metadata_file_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_metadata_file_ids),
icda_events
AS (INSERT INTO cleanup_event_icda_association
(create_time, cleanup_event_id, icda_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_icda_ids),
icda_hda_events
AS (INSERT INTO cleanup_event_hda_association
(create_time, cleanup_event_id, hda_id)
SELECT NOW(), %s, hda_id
SELECT NOW(), %(event_id)s, hda_id
FROM deleted_icda_ids)
SELECT purged_hda_ids.id,
history.user_id,
@@ -438,24 +463,10 @@ class Cleanup(object):
LEFT OUTER JOIN deleted_icda_ids
ON deleted_icda_ids.hda_parent_id = purged_hda_ids.id
LEFT OUTER JOIN history
ON purged_hda_ids.history_id = history.id;
ON purged_hda_ids.history_id = history.id
"""
force_retry_sql = """
AND NOT purged"""
update_time_sql = ""
if self.args.force_retry:
force_retry_sql = ""
else:
# only update time if not doing force retry (otherwise a lot of things would have their update times reset that were actually purged a long time ago)
if self.args.update_time:
update_time_sql = """,
update_time = NOW()"""
sql = sql % (update_time_sql, force_retry_sql, '%s', update_time_sql, update_time_sql, update_time_sql, '%s', '%s', '%s', '%s')
args = (self.args.days, event_id, event_id, event_id, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -477,19 +488,17 @@ class Cleanup(object):
"""
log.info('Marking purged all deleted histories that are older than the specified number of days.')
event_id = self._create_event()
sql = """
WITH purged_history_ids
AS ( UPDATE history
SET purged = true%s
WHERE deleted%s
AND update_time < (NOW() - interval '%s days')
SET purged = true{update_time_sql}
WHERE deleted{force_retry_sql}
AND update_time < (NOW() - interval '%(days)s days')
RETURNING id,
user_id),
purged_hda_ids
AS ( UPDATE history_dataset_association
SET purged = true%s
SET purged = true{update_time_sql}
FROM purged_history_ids
WHERE purged_history_ids.id = history_dataset_association.history_id
AND NOT history_dataset_association.purged
@@ -497,7 +506,7 @@ class Cleanup(object):
history_dataset_association.id AS id),
deleted_metadata_file_ids
AS ( UPDATE metadata_file
SET deleted = true%s
SET deleted = true{update_time_sql}
FROM purged_hda_ids
WHERE purged_hda_ids.id = metadata_file.hda_id
RETURNING metadata_file.hda_id AS hda_id,
@@ -505,7 +514,7 @@ class Cleanup(object):
metadata_file.object_store_id AS object_store_id),
deleted_icda_ids
AS ( UPDATE implicitly_converted_dataset_association
SET deleted = true%s
SET deleted = true{update_time_sql}
FROM purged_hda_ids
WHERE purged_hda_ids.id = implicitly_converted_dataset_association.hda_parent_id
RETURNING implicitly_converted_dataset_association.hda_id AS hda_id,
@@ -513,33 +522,33 @@ class Cleanup(object):
implicitly_converted_dataset_association.id AS id),
deleted_icda_purged_child_hda_ids
AS ( UPDATE history_dataset_association
SET purged = true%s
SET purged = true{update_time_sql}
FROM deleted_icda_ids
WHERE deleted_icda_ids.hda_id = history_dataset_association.id),
history_events
AS (INSERT INTO cleanup_event_history_association
(create_time, cleanup_event_id, history_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM purged_history_ids),
hda_events
AS (INSERT INTO cleanup_event_hda_association
(create_time, cleanup_event_id, hda_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM purged_hda_ids),
metadata_file_events
AS (INSERT INTO cleanup_event_metadata_file_association
(create_time, cleanup_event_id, metadata_file_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_metadata_file_ids),
icda_events
AS (INSERT INTO cleanup_event_icda_association
(create_time, cleanup_event_id, icda_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_icda_ids),
icda_hda_events
AS (INSERT INTO cleanup_event_hda_association
(create_time, cleanup_event_id, hda_id)
SELECT NOW(), %s, hda_id
SELECT NOW(), %(event_id)s, hda_id
FROM deleted_icda_ids)
SELECT purged_history_ids.id,
purged_history_ids.user_id,
@@ -554,23 +563,11 @@ class Cleanup(object):
LEFT OUTER JOIN deleted_metadata_file_ids
ON deleted_metadata_file_ids.hda_id = purged_hda_ids.id
LEFT OUTER JOIN deleted_icda_ids
ON deleted_icda_ids.hda_parent_id = purged_hda_ids.id;
ON deleted_icda_ids.hda_parent_id = purged_hda_ids.id
ORDER BY purged_history_ids.id
"""
force_retry_sql = """
AND NOT purged"""
update_time_sql = ""
if self.args.force_retry:
force_retry_sql = ""
else:
if self.args.update_time:
update_time_sql += """,
update_time = NOW()"""
sql = sql % (update_time_sql, force_retry_sql, '%s', update_time_sql, update_time_sql, update_time_sql, update_time_sql, '%s', '%s', '%s', '%s', '%s')
args = (self.args.days, event_id, event_id, event_id, event_id, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -593,35 +590,26 @@ class Cleanup(object):
"""
log.info('Marking deleted all Datasets that are derivative of JobExportHistoryArchives that are older than the specified number of days.')
event_id = self._create_event()
sql = """
WITH deleted_dataset_ids
AS ( UPDATE dataset
SET deleted = true%s
SET deleted = true{update_time_sql}
FROM job_export_history_archive
WHERE job_export_history_archive.dataset_id = dataset.id
AND NOT deleted
AND dataset.update_time <= (NOW() - interval '%s days')
AND dataset.update_time <= (NOW() - interval '%(days)s days')
RETURNING dataset.id),
dataset_events
AS (INSERT INTO cleanup_event_dataset_association
(create_time, cleanup_event_id, dataset_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_dataset_ids)
SELECT id
FROM deleted_dataset_ids
ORDER BY id;
ORDER BY id
"""
update_time_sql = ""
if self.args.update_time:
update_time_sql += """,
update_time = NOW()"""
sql = sql % (update_time_sql, '%s', '%s')
args = (self.args.days, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -635,42 +623,33 @@ class Cleanup(object):
"""
log.info('Marking deleted all Datasets whose associations are all marked as deleted/purged that are older than the specified number of days.')
event_id = self._create_event()
sql = """
WITH deleted_dataset_ids
AS ( UPDATE dataset
SET deleted = true%s
SET deleted = true{update_time_sql}
WHERE NOT deleted
AND NOT EXISTS (SELECT true
FROM library_dataset_dataset_association
WHERE (NOT deleted
OR update_time >= (NOW() - interval '%s days'))
OR update_time >= (NOW() - interval '%(days)s days'))
AND dataset.id = dataset_id)
AND NOT EXISTS (SELECT true
FROM history_dataset_association
WHERE (NOT purged
OR update_time >= (NOW() - interval '%s days'))
OR update_time >= (NOW() - interval '%(days)s days'))
AND dataset.id = dataset_id)
RETURNING id),
dataset_events
AS (INSERT INTO cleanup_event_dataset_association
(create_time, cleanup_event_id, dataset_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM deleted_dataset_ids)
SELECT id
FROM deleted_dataset_ids
ORDER BY id;
ORDER BY id
"""
update_time_sql = ""
if self.args.update_time:
update_time_sql += """,
update_time = NOW()"""
sql = sql % (update_time_sql, '%s', '%s', '%s')
args = (self.args.days, self.args.days, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -684,41 +663,26 @@ class Cleanup(object):
"""
log.info('Marking purged all Datasets marked deleted that are older than the specified number of days.')
event_id = self._create_event()
sql = """
WITH purged_dataset_ids
AS ( UPDATE dataset
SET purged = true%s
WHERE deleted%s
AND update_time < (NOW() - interval '%s days')
SET purged = true{update_time_sql}
WHERE deleted{force_retry_sql}
AND update_time < (NOW() - interval '%(days)s days')
RETURNING id,
object_store_id),
dataset_events
AS (INSERT INTO cleanup_event_dataset_association
(create_time, cleanup_event_id, dataset_id)
SELECT NOW(), %s, id
SELECT NOW(), %(event_id)s, id
FROM purged_dataset_ids)
SELECT id,
object_store_id
FROM purged_dataset_ids
ORDER BY id;
ORDER BY id
"""
force_retry_sql = """
AND NOT purged"""
update_time_sql = ""
if self.args.force_retry:
force_retry_sql = ""
else:
if self.args.update_time:
update_time_sql = """,
update_time = NOW()"""
sql = sql % (update_time_sql, force_retry_sql, '%s', '%s')
args = (self.args.days, event_id)
cur = self._update(sql, args)
cur = self._update(sql, {})
self._flush()
self._open_logfile()
@@ -758,6 +722,29 @@ class Cleanup(object):
self._close_logfile()
def quote_for_print(v):
if isinstance(v, string_types):
return "'" + v.replace("'", "\\'") + "'"
else:
return v
def args_for_print(args):
if isinstance(args, dict):
r = {}
for k, v in args.items():
r[k] = quote_for_print(v)
elif isinstance(args, (tuple, list)):
r = []
for v in args:
r.append(quote_for_print(v))
if isinstance(args, tuple):
r = tuple(r)
else:
r = args
return r
if __name__ == '__main__':
cleanup = Cleanup()
try: