diff --git a/rowers/longtask.py b/rowers/longtask.py new file mode 100644 index 00000000..e0033967 --- /dev/null +++ b/rowers/longtask.py @@ -0,0 +1,58 @@ +from __future__ import absolute_import +import numpy as np +import time + +from redis import StrictRedis,Redis +from celery import result as celery_result +import json + +redis_connection = StrictRedis() + +import redis +import threading + +def getvalue(data): + perc = 0 + total = 1 + done = 0 + id = 0 + session_key = 'noot' + for i in data.iteritems(): + if i[0] == 'total': + total = float(i[1]) + if i[0] == 'done': + done = float(i[1]) + if i[0] == 'id': + id = i[1] + if i[0] == 'session_key': + session_key = i[1] + + return total,done,id,session_key + + + +def longtask(aantal,jobid=None,debug=False, + session_key=None): + counter = 0 + + channel = 'tasks' + for i in range(aantal): + time.sleep(1) + counter += 1 + if counter > 10: + counter = 0 + if debug: + progress = 100.*i/aantal + if jobid != None: + redis_connection.publish(channel,json.dumps( + { + 'done':i, + 'total':aantal, + 'id':jobid, + 'session_key':session_key, + } + )) + + + + return 1 diff --git a/rowers/tasks.py b/rowers/tasks.py index 237f9be4..ae656006 100644 --- a/rowers/tasks.py +++ b/rowers/tasks.py @@ -35,6 +35,8 @@ from django.db.utils import OperationalError import datautils import utils +import longtask + # testing task @@ -42,9 +44,15 @@ import utils def add(x, y): return x + y + +@app.task(bind=True) +def long_test_task(self,aantal,debug=False,job=None,session_key=None): + job = self.request + + return longtask.longtask(aantal,jobid=job.id,debug=debug, + session_key=session_key) + # create workout - - @app.task def handle_new_workout_from_file(r, f2, workouttype='rower', diff --git a/rowers/templates/async_tasks.html b/rowers/templates/async_tasks.html index 1986acfb..dedf883d 100644 --- a/rowers/templates/async_tasks.html +++ b/rowers/templates/async_tasks.html @@ -4,6 +4,53 @@ {% block title %}Rowsandall - Tasks {% endblock %} +{% block meta %} + + + + + +{% endblock %} + {% block content %} @@ -16,7 +63,8 @@ ID - Task + Task + Progress Status Action @@ -31,6 +79,12 @@ {{ task|lookup:'verbose' }} + + {{ task|lookup:'progress' }} + + + + {{ task|lookup:'status' }} {% if task|lookup:'failed' %} diff --git a/rowers/templatetags/rowerfilters.py b/rowers/templatetags/rowerfilters.py index 123faa9e..eebaaaf3 100644 --- a/rowers/templatetags/rowerfilters.py +++ b/rowers/templatetags/rowerfilters.py @@ -1,6 +1,8 @@ from django import template +from django.utils.safestring import mark_safe from time import strftime import dateutil.parser +import json register = template.Library() @@ -65,6 +67,11 @@ def deltatimeprint(d): return strfdeltah(d) +@register.filter(is_safe=True) +def jsdict(dict,key): + s = dict.get(key) + return mark_safe(json.dumps(s)) + @register.filter def lookup(dict, key): s = dict.get(key) diff --git a/rowers/urls.py b/rowers/urls.py index 84a78222..f592717a 100644 --- a/rowers/urls.py +++ b/rowers/urls.py @@ -141,6 +141,7 @@ urlpatterns = [ url(r'^list-jobs/$',views.session_jobs_view), url(r'^jobs-status/$',views.session_jobs_status), url(r'^job-kill/(?P.*)$',views.kill_async_job), + url(r'^test-job/(?P\d+)$',views.test_job_view), url(r'^list-graphs/$',views.graphs_view), url(r'^(?P\d+)/ote-bests/(?P\w+.*)/(?P\w+.*)$',views.rankings_view), url(r'^(?P\d+)/ote-bests/(?P\d+)$',views.rankings_view), diff --git a/rowers/utils.py b/rowers/utils.py index b6b81eec..03c59140 100644 --- a/rowers/utils.py +++ b/rowers/utils.py @@ -5,6 +5,7 @@ import colorsys from django.conf import settings + lbstoN = 4.44822 landingpages = ( @@ -226,3 +227,5 @@ def myqueue(queue,function,*args,**kwargs): job = queue.enqueue(function,*args,**kwargs) return job + + diff --git a/rowers/views.py b/rowers/views.py index 0ad77a0f..7f0729ec 100644 --- a/rowers/views.py +++ b/rowers/views.py @@ -14,6 +14,7 @@ from django.views.generic.base import TemplateView from django.db.models import Q from django import template from django.db import IntegrityError, transaction + from django.shortcuts import render from django.http import ( HttpResponse, HttpResponseRedirect, @@ -99,7 +100,7 @@ from rowers.tasks import handle_makeplot,handle_otwsetpower,handle_sendemailtcx, from rowers.tasks import ( handle_sendemail_unrecognized,handle_sendemailnewcomment, handle_sendemailnewresponse, handle_updatedps, - handle_updatecp + handle_updatecp,long_test_task ) from scipy.signal import savgol_filter @@ -130,14 +131,70 @@ import django_rq queue = django_rq.get_queue('default') queuelow = django_rq.get_queue('low') queuehigh = django_rq.get_queue('low') - + +import redis +import threading from redis import StrictRedis,Redis from rq.exceptions import NoSuchJobError from rq.registry import StartedJobRegistry from rq import Queue,cancel_job +from django.core.cache import cache + +def getvalue(data): + perc = 0 + total = 1 + done = 0 + id = 0 + session_key = 'noot' + for i in data.iteritems(): + if i[0] == 'total': + total = float(i[1]) + if i[0] == 'done': + done = float(i[1]) + if i[0] == 'id': + id = i[1] + if i[0] == 'session_key': + session_key = i[1] + + return total,done,id,session_key + +class SessionTaskListener(threading.Thread): + def __init__(self, r, channels): + threading.Thread.__init__(self) + self.redis = r + self.pubsub = self.redis.pubsub() + self.pubsub.subscribe(channels) + + def work(self, item): + + try: + data = json.loads(item['data']) + total,done,id,session_key = getvalue(data) + perc = int(100.*done/total) + cache.set(id,perc,3600) + + except TypeError: + pass + + def run(self): + for item in self.pubsub.listen(): + if item['data'] == "KILL": + self.pubsub.unsubscribe() + print self, "unsubscribed and finished" + break + else: + self.work(item) + + queuefailed = Queue("failed",connection=Redis()) redis_connection = StrictRedis() +r = Redis() + +client = SessionTaskListener(r,['tasks']) +client.start() + + rq_registry = StartedJobRegistry(queue.name,connection=redis_connection) rq_registryhigh = StartedJobRegistry(queuehigh.name,connection=redis_connection) rq_registrylow = StartedJobRegistry(queuelow.name,connection=redis_connection) @@ -181,14 +238,13 @@ def remove_asynctask(request,id): newtasks = [] for task in oldtasks: - print task[0] if id not in task[0]: newtasks += [(task[0],task[1])] request.session['async_tasks'] = newtasks def get_job_result(jobid): - if settings.EBUG: + if settings.DEBUG: result = celery_result.AsyncResult(jobid).result else: running_job_ids = rq_registry.get_job_ids() @@ -210,12 +266,14 @@ verbose_job_status = { 'updatecpwater': 'Critical Power Calculation for OTW Workouts', 'otwsetpower': 'Rowing Physics OTW Power Calculation', 'make_plot': 'Create static chart', + 'long_test_task': 'Long Test Task', } def get_job_status(jobid): if settings.DEBUG: job = celery_result.AsyncResult(jobid) jobresult = job.result + if 'fail' in job.status.lower(): jobresult = '0' summary = { @@ -261,10 +319,29 @@ def kill_async_job(request,id='aap'): pass remove_asynctask(request,id) + cache.delete(id) url = reverse(session_jobs_status) return HttpResponseRedirect(url) - + +@login_required() +def test_job_view(request,aantal=100): + + session_key = request.session._session_key + + job = myqueue(queuehigh,long_test_task,int(aantal), + session_key=session_key) + + + try: + request.session['async_tasks'] += [(job.id,'long_test_task')] + except KeyError: + request.session['async_tasks'] = [(job.id,'long_test_task')] + + url = reverse(session_jobs_status) + + return HttpResponseRedirect(url) + def get_all_queued_jobs(userid=0): r = StrictRedis() @@ -310,15 +387,32 @@ def get_stored_tasks_status(request): except KeyError: taskids = [] - taskstatus = [{ - 'id':id, - 'status':get_job_status(id)['status'], - 'failed':get_job_status(id)['failed'], - 'finished':get_job_status(id)['finished'], - 'func_name':func_name, - 'verbose': verbose_job_status[func_name] - } for id,func_name in taskids] + taskstatus = [] + for id,func_name in taskids: + progress = 0 + cached_progress = cache.get(id) + finished = get_job_status(id)['finished'] + if finished: + cache.set(id,100) + progress = 100 + elif cached_progress>0: + progress = cached_progress + else: + progress = 0 + this_task_status = { + 'id':id, + 'status':get_job_status(id)['status'], + 'failed':get_job_status(id)['failed'], + 'finished':get_job_status(id)['finished'], + 'func_name':func_name, + 'verbose': verbose_job_status[func_name], + 'progress': progress, + } + + taskstatus.append(this_task_status) + + return taskstatus @login_required() @@ -6528,7 +6622,6 @@ def workout_flexchart3_view(request,*args,**kwargs): if request.user == row.user.user: mayedit=1 - workouttype = 'ote' if row.workouttype in ('water','coastal'): workouttype = 'otw' @@ -6578,8 +6671,9 @@ def workout_flexchart3_view(request,*args,**kwargs): else: yparam2 = 'hr' - if favoritenr>=0 and r.showfavoritechartnotes: - favoritechartnotes = favorites[favoritenr].notes + if not request.user.is_anonymous(): + if favoritenr>=0 and r.showfavoritechartnotes: + favoritechartnotes = favorites[favoritenr].notes else: favoritechartnotes = '' favoritenr = 0 diff --git a/rowsandall_app/settings.py b/rowsandall_app/settings.py index beadcf86..7f8a63b2 100644 --- a/rowsandall_app/settings.py +++ b/rowsandall_app/settings.py @@ -287,6 +287,7 @@ RQ_QUEUES = { #SESSION_ENGINE = "django.contrib.sessions.backends.signed_cookies" #SESSION_ENGINE = "django.contrib.sessions.backends.cached_db" SESSION_ENGINE = "django.contrib.sessions.backends.cache" +SESSION_SAVE_EVERY_REQUEST = True # admin stuff for error reporting SERVER_EMAIL='admin@rowsandall.com'