From f2db4b45dde3a4575b17af9b968d09820273db82 Mon Sep 17 00:00:00 2001 From: Pierre-Yves Chibon Date: Dec 18 2017 09:57:21 +0000 Subject: [PATCH 1/4] Bind the tasks so they can update their status and report progress Signed-off-by: Pierre-Yves Chibon --- diff --git a/pagure/lib/tasks.py b/pagure/lib/tasks.py index 337c6a1..96e27ac 100644 --- a/pagure/lib/tasks.py +++ b/pagure/lib/tasks.py @@ -74,8 +74,8 @@ def gc_clean(): gc.collect() -@conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None)) -def generate_gitolite_acls(namespace=None, name=None, user=None, group=None): +@conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None), bind=True) +def generate_gitolite_acls(self, namespace=None, name=None, user=None, group=None): """ Generate the gitolite configuration file either entirely or for a specific project. @@ -89,6 +89,9 @@ def generate_gitolite_acls(namespace=None, name=None, user=None, group=None): :type group: None or str """ + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = None if name and name != -1: @@ -122,8 +125,8 @@ def generate_gitolite_acls(namespace=None, name=None, user=None, group=None): gc_clean() -@conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None)) -def delete_project(namespace=None, name=None, user=None, action_user=None): +@conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None), bind=True) +def delete_project(self, namespace=None, name=None, user=None, action_user=None): """ Delete a project in pagure. This is achieved in three steps: @@ -141,6 +144,9 @@ def delete_project(namespace=None, name=None, user=None, action_user=None): :type action_user: None or str """ + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( session, namespace=namespace, name=name, user=user, @@ -207,8 +213,8 @@ def delete_project(namespace=None, name=None, user=None, action_user=None): return ret('view_user', username=username) -@conn.task -def create_project(username, namespace, name, add_readme, +@conn.task(bind=True) +def create_project(self, username, namespace, name, add_readme, ignore_existing_repo): """ Create a project. @@ -226,6 +232,9 @@ def create_project(username, namespace, name, add_readme, :type ignore_existing_repo: bool """ + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -333,8 +342,11 @@ def create_project(username, namespace, name, add_readme, return ret('view_repo', repo=name, namespace=namespace) -@conn.task -def update_git(name, namespace, user, ticketuid=None, requestuid=None): +@conn.task(bind=True) +def update_git(self, name, namespace, user, ticketuid=None, requestuid=None): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -361,8 +373,11 @@ def update_git(name, namespace, user, ticketuid=None, requestuid=None): return result -@conn.task -def clean_git(name, namespace, user, ticketuid): +@conn.task(bind=True) +def clean_git(self, name, namespace, user, ticketuid): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -382,9 +397,14 @@ def clean_git(name, namespace, user, ticketuid): return result -@conn.task -def update_file_in_git(name, namespace, user, branch, branchto, filename, - content, message, username, email, runhook=False): +@conn.task(bind=True) +def update_file_in_git(self, name, namespace, user, branch, branchto, + filename, content, message, username, email, + runhook=False): + + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() userobj = pagure.lib.search_user(session, username=username) @@ -402,8 +422,11 @@ def update_file_in_git(name, namespace, user, branch, branchto, filename, namespace=namespace, branchname=branchto) -@conn.task -def delete_branch(name, namespace, user, branchname): +@conn.task(bind=True) +def delete_branch(self, name, namespace, user, branchname): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -423,8 +446,8 @@ def delete_branch(name, namespace, user, branchname): return ret('view_repo', repo=name, namespace=namespace, username=user) -@conn.task -def fork(name, namespace, user_owner, user_forker, editbranch, editfile): +@conn.task(bind=True) +def fork(self, name, namespace, user_owner, user_forker, editbranch, editfile): """ Forks the specified project for the specified user. :arg namespace: the namespace of the project @@ -443,6 +466,9 @@ def fork(name, namespace, user_owner, user_forker, editbranch, editfile): :type editfile: str """ + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() repo_from = pagure.lib._get_project( @@ -538,8 +564,11 @@ def fork(name, namespace, user_owner, user_forker, editbranch, editfile): filename=editfile) -@conn.task -def pull_remote_repo(remote_git, branch_from): +@conn.task(bind=True) +def pull_remote_repo(self, remote_git, branch_from): + if self is not None: + self.update_state(state='RUNNING') + clonepath = pagure.get_remote_repo_path(remote_git, branch_from, ignore_non_exist=True) repo = pygit2.clone_repository( @@ -550,8 +579,11 @@ def pull_remote_repo(remote_git, branch_from): return clonepath -@conn.task -def refresh_remote_pr(name, namespace, user, requestid): +@conn.task(bind=True) +def refresh_remote_pr(self, name, namespace, user, requestid): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -578,8 +610,11 @@ def refresh_remote_pr(name, namespace, user, requestid): requestid=requestid) -@conn.task -def refresh_pr_cache(name, namespace, user): +@conn.task(bind=True) +def refresh_pr_cache(self, name, namespace, user): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -592,8 +627,11 @@ def refresh_pr_cache(name, namespace, user): gc_clean() -@conn.task -def merge_pull_request(name, namespace, user, requestid, user_merger): +@conn.task(bind=True) +def merge_pull_request(self, name, namespace, user, requestid, user_merger): + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -615,8 +653,11 @@ def merge_pull_request(name, namespace, user, requestid, user_merger): return ret('view_repo', repo=name, username=user, namespace=namespace) -@conn.task -def add_file_to_git(name, namespace, user, user_attacher, issueuid, filename): +@conn.task(bind=True) +def add_file_to_git( + self, name, namespace, user, user_attacher, issueuid, filename): + if self is not None: + self.update_state(state='RUNNING') session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -642,13 +683,16 @@ def add_file_to_git(name, namespace, user, user_attacher, issueuid, filename): gc_clean() -@conn.task -def project_dowait(name, namespace, user): +@conn.task(bind=True) +def project_dowait(self, name, namespace, user): """ This is a task used to test the locking systems. It should never be allowed to be called in production instances, since that would allow an attacker to basically DOS a project by calling this repeatedly. """ + if self is not None: + self.update_state(state='RUNNING') + assert APP.config.get('ALLOW_PROJECT_DOWAIT', False) session = pagure.lib.create_session() @@ -666,11 +710,14 @@ def project_dowait(name, namespace, user): return ret('view_repo', repo=name, username=user, namespace=namespace) -@conn.task -def sync_pull_ref(name, namespace, user, requestid): +@conn.task(bind=True) +def sync_pull_ref(self, name, namespace, user, requestid): """ Synchronize a pull/ reference from the content in the forked repo, allowing local checkout of the pull-request. """ + if self is not None: + self.update_state(state='RUNNING') + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -700,10 +747,13 @@ def sync_pull_ref(name, namespace, user, requestid): gc_clean() -@conn.task -def update_checksums_file(folder, filenames): +@conn.task(bind=True) +def update_checksums_file(self, folder, filenames): """ """ + if self is not None: + self.update_state(state='RUNNING') + sha_file = os.path.join(folder, 'CHECKSUMS') new_file = not os.path.exists(sha_file) @@ -739,11 +789,14 @@ def update_checksums_file(folder, filenames): algo.upper(), filename, algos[algo].hexdigest())) -@conn.task -def commits_author_stats(repopath): +@conn.task(bind=True) +def commits_author_stats(self, repopath): """ Returns some statistics about commits made against the specified git repository. """ + if self is not None: + self.update_state(state='RUNNING') + if not os.path.exists(repopath): raise ValueError('Git repository not found.') @@ -771,11 +824,14 @@ def commits_author_stats(repopath): return (cnt, out_list, len(authors_email), commit.commit_time) -@conn.task -def commits_history_stats(repopath): +@conn.task(bind=True) +def commits_history_stats(self, repopath): """ Returns the evolution of the commits made against the specified git repository. """ + if self is not None: + self.update_state(state='RUNNING') + if not os.path.exists(repopath): raise ValueError('Git repository not found.') From c2199b3ff13f3d9430d0edb7cca237168a125dd1 Mon Sep 17 00:00:00 2001 From: Pierre-Yves Chibon Date: Dec 18 2017 10:20:21 +0000 Subject: [PATCH 2/4] Re-order the imports Signed-off-by: Pierre-Yves Chibon --- diff --git a/pagure/lib/tasks.py b/pagure/lib/tasks.py index 96e27ac..66c28f0 100644 --- a/pagure/lib/tasks.py +++ b/pagure/lib/tasks.py @@ -12,21 +12,21 @@ import collections import datetime import gc import hashlib +import logging import os import os.path import shutil +import tempfile import time -from celery import Celery -from celery.result import AsyncResult +from functools import wraps import arrow import pygit2 -import tempfile import six -import logging - +from celery import Celery +from celery.result import AsyncResult from sqlalchemy.exc import SQLAlchemyError import pagure From bbfb815ad756d795be1a38ed0d88567a6f7caa79 Mon Sep 17 00:00:00 2001 From: Pierre-Yves Chibon Date: Dec 18 2017 10:21:14 +0000 Subject: [PATCH 3/4] Move to a decorator instead of re-using the same code everywhere Signed-off-by: Pierre-Yves Chibon --- diff --git a/pagure/lib/tasks.py b/pagure/lib/tasks.py index 66c28f0..317fc1d 100644 --- a/pagure/lib/tasks.py +++ b/pagure/lib/tasks.py @@ -51,6 +51,19 @@ conn = Celery('tasks', broker=broker_url, backend=broker_url) conn.conf.update(APP.config['CELERY_CONFIG']) +def set_status(function): + """ Simple decorator adjusting the status of the task when it starts. + """ + + @wraps(function) + def decorated_function(self, *args, **kwargs): + """ Decorated function, actually does the work. """ + if self is not None: + self.update_state(state='RUNNING') + return function(self, *args, **kwargs) + return decorated_function + + def get_result(uuid): """ Returns the AsyncResult object for a given task. @@ -75,6 +88,7 @@ def gc_clean(): @conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None), bind=True) +@set_status def generate_gitolite_acls(self, namespace=None, name=None, user=None, group=None): """ Generate the gitolite configuration file either entirely or for a specific project. @@ -89,8 +103,6 @@ def generate_gitolite_acls(self, namespace=None, name=None, user=None, group=Non :type group: None or str """ - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() project = None @@ -126,6 +138,7 @@ def generate_gitolite_acls(self, namespace=None, name=None, user=None, group=Non @conn.task(queue=APP.config.get('GITOLITE_CELERY_QUEUE', None), bind=True) +@set_status def delete_project(self, namespace=None, name=None, user=None, action_user=None): """ Delete a project in pagure. @@ -144,8 +157,6 @@ def delete_project(self, namespace=None, name=None, user=None, action_user=None) :type action_user: None or str """ - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -214,6 +225,7 @@ def delete_project(self, namespace=None, name=None, user=None, action_user=None) @conn.task(bind=True) +@set_status def create_project(self, username, namespace, name, add_readme, ignore_existing_repo): """ Create a project. @@ -232,8 +244,6 @@ def create_project(self, username, namespace, name, add_readme, :type ignore_existing_repo: bool """ - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -343,9 +353,8 @@ def create_project(self, username, namespace, name, add_readme, @conn.task(bind=True) +@set_status def update_git(self, name, namespace, user, ticketuid=None, requestuid=None): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -374,9 +383,8 @@ def update_git(self, name, namespace, user, ticketuid=None, requestuid=None): @conn.task(bind=True) +@set_status def clean_git(self, name, namespace, user, ticketuid): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -398,13 +406,11 @@ def clean_git(self, name, namespace, user, ticketuid): @conn.task(bind=True) +@set_status def update_file_in_git(self, name, namespace, user, branch, branchto, filename, content, message, username, email, runhook=False): - if self is not None: - self.update_state(state='RUNNING') - session = pagure.lib.create_session() userobj = pagure.lib.search_user(session, username=username) @@ -423,9 +429,8 @@ def update_file_in_git(self, name, namespace, user, branch, branchto, @conn.task(bind=True) +@set_status def delete_branch(self, name, namespace, user, branchname): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -447,6 +452,7 @@ def delete_branch(self, name, namespace, user, branchname): @conn.task(bind=True) +@set_status def fork(self, name, namespace, user_owner, user_forker, editbranch, editfile): """ Forks the specified project for the specified user. @@ -466,8 +472,6 @@ def fork(self, name, namespace, user_owner, user_forker, editbranch, editfile): :type editfile: str """ - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -565,9 +569,8 @@ def fork(self, name, namespace, user_owner, user_forker, editbranch, editfile): @conn.task(bind=True) +@set_status def pull_remote_repo(self, remote_git, branch_from): - if self is not None: - self.update_state(state='RUNNING') clonepath = pagure.get_remote_repo_path(remote_git, branch_from, ignore_non_exist=True) @@ -580,9 +583,8 @@ def pull_remote_repo(self, remote_git, branch_from): @conn.task(bind=True) +@set_status def refresh_remote_pr(self, name, namespace, user, requestid): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -611,9 +613,8 @@ def refresh_remote_pr(self, name, namespace, user, requestid): @conn.task(bind=True) +@set_status def refresh_pr_cache(self, name, namespace, user): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -628,9 +629,8 @@ def refresh_pr_cache(self, name, namespace, user): @conn.task(bind=True) +@set_status def merge_pull_request(self, name, namespace, user, requestid, user_merger): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -654,10 +654,9 @@ def merge_pull_request(self, name, namespace, user, requestid, user_merger): @conn.task(bind=True) +@set_status def add_file_to_git( self, name, namespace, user, user_attacher, issueuid, filename): - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -684,14 +683,13 @@ def add_file_to_git( @conn.task(bind=True) +@set_status def project_dowait(self, name, namespace, user): """ This is a task used to test the locking systems. It should never be allowed to be called in production instances, since that would allow an attacker to basically DOS a project by calling this repeatedly. """ - if self is not None: - self.update_state(state='RUNNING') assert APP.config.get('ALLOW_PROJECT_DOWAIT', False) @@ -711,12 +709,11 @@ def project_dowait(self, name, namespace, user): @conn.task(bind=True) +@set_status def sync_pull_ref(self, name, namespace, user, requestid): """ Synchronize a pull/ reference from the content in the forked repo, allowing local checkout of the pull-request. """ - if self is not None: - self.update_state(state='RUNNING') session = pagure.lib.create_session() @@ -748,11 +745,10 @@ def sync_pull_ref(self, name, namespace, user, requestid): @conn.task(bind=True) +@set_status def update_checksums_file(self, folder, filenames): """ """ - if self is not None: - self.update_state(state='RUNNING') sha_file = os.path.join(folder, 'CHECKSUMS') new_file = not os.path.exists(sha_file) @@ -790,12 +786,11 @@ def update_checksums_file(self, folder, filenames): @conn.task(bind=True) +@set_status def commits_author_stats(self, repopath): """ Returns some statistics about commits made against the specified git repository. """ - if self is not None: - self.update_state(state='RUNNING') if not os.path.exists(repopath): raise ValueError('Git repository not found.') @@ -825,12 +820,11 @@ def commits_author_stats(self, repopath): @conn.task(bind=True) +@set_status def commits_history_stats(self, repopath): """ Returns the evolution of the commits made against the specified git repository. """ - if self is not None: - self.update_state(state='RUNNING') if not os.path.exists(repopath): raise ValueError('Git repository not found.') From 314da360fc8aa449df663d6f86b25553bae09f49 Mon Sep 17 00:00:00 2001 From: Pierre-Yves Chibon Date: Dec 18 2017 10:21:47 +0000 Subject: [PATCH 4/4] Add docstrings to the functions missing one Signed-off-by: Pierre-Yves Chibon --- diff --git a/pagure/lib/tasks.py b/pagure/lib/tasks.py index 317fc1d..6a170cd 100644 --- a/pagure/lib/tasks.py +++ b/pagure/lib/tasks.py @@ -355,6 +355,9 @@ def create_project(self, username, namespace, name, add_readme, @conn.task(bind=True) @set_status def update_git(self, name, namespace, user, ticketuid=None, requestuid=None): + """ Update the JSON representation of either a ticket or a pull-request + depending on the argument specified. + """ session = pagure.lib.create_session() @@ -385,6 +388,9 @@ def update_git(self, name, namespace, user, ticketuid=None, requestuid=None): @conn.task(bind=True) @set_status def clean_git(self, name, namespace, user, ticketuid): + """ Remove the JSON representation of a ticket on the git repository + for tickets. + """ session = pagure.lib.create_session() @@ -410,6 +416,8 @@ def clean_git(self, name, namespace, user, ticketuid): def update_file_in_git(self, name, namespace, user, branch, branchto, filename, content, message, username, email, runhook=False): + """ Update a file in the specified git repo. + """ session = pagure.lib.create_session() @@ -431,6 +439,8 @@ def update_file_in_git(self, name, namespace, user, branch, branchto, @conn.task(bind=True) @set_status def delete_branch(self, name, namespace, user, branchname): + """ Delete a branch from a git repo. + """ session = pagure.lib.create_session() @@ -571,6 +581,8 @@ def fork(self, name, namespace, user_owner, user_forker, editbranch, editfile): @conn.task(bind=True) @set_status def pull_remote_repo(self, remote_git, branch_from): + """ Clone a remote git repository locally for remote PRs. + """ clonepath = pagure.get_remote_repo_path(remote_git, branch_from, ignore_non_exist=True) @@ -585,6 +597,9 @@ def pull_remote_repo(self, remote_git, branch_from): @conn.task(bind=True) @set_status def refresh_remote_pr(self, name, namespace, user, requestid): + """ Refresh the local clone of a git repository used in a remote + pull-request. + """ session = pagure.lib.create_session() @@ -615,6 +630,8 @@ def refresh_remote_pr(self, name, namespace, user, requestid): @conn.task(bind=True) @set_status def refresh_pr_cache(self, name, namespace, user): + """ Refresh the merge status cached of pull-requests. + """ session = pagure.lib.create_session() @@ -631,6 +648,8 @@ def refresh_pr_cache(self, name, namespace, user): @conn.task(bind=True) @set_status def merge_pull_request(self, name, namespace, user, requestid, user_merger): + """ Merge pull-request. + """ session = pagure.lib.create_session() @@ -657,6 +676,9 @@ def merge_pull_request(self, name, namespace, user, requestid, user_merger): @set_status def add_file_to_git( self, name, namespace, user, user_attacher, issueuid, filename): + """ Add a file to the specified git repo. + """ + session = pagure.lib.create_session() project = pagure.lib._get_project( @@ -747,7 +769,7 @@ def sync_pull_ref(self, name, namespace, user, requestid): @conn.task(bind=True) @set_status def update_checksums_file(self, folder, filenames): - """ + """ Update the checksums file in the release folder of the project. """ sha_file = os.path.join(folder, 'CHECKSUMS')