|
@@ -2,10 +2,12 @@ import os.path
|
|
|
import datetime
|
|
import datetime
|
|
|
import json
|
|
import json
|
|
|
import thumb
|
|
import thumb
|
|
|
|
|
+import time
|
|
|
|
|
|
|
|
from constants import CUT_OUTPUT_PATH
|
|
from constants import CUT_OUTPUT_PATH
|
|
|
from subprocess import Popen,PIPE
|
|
from subprocess import Popen,PIPE
|
|
|
from rq import get_current_job,Queue
|
|
from rq import get_current_job,Queue
|
|
|
|
|
+from rq.registry import StartedJobRegistry,FinishedJobRegistry
|
|
|
from rq.job import Job
|
|
from rq.job import Job
|
|
|
from rq.decorators import job
|
|
from rq.decorators import job
|
|
|
from worker import conn,q
|
|
from worker import conn,q
|
|
@@ -13,17 +15,8 @@ from dateutil import tz
|
|
|
from db import db
|
|
from db import db
|
|
|
|
|
|
|
|
d = db()
|
|
d = db()
|
|
|
-job_key_list = []
|
|
|
|
|
date_fmt='%d/%m/%Y, %H:%M'
|
|
date_fmt='%d/%m/%Y, %H:%M'
|
|
|
|
|
|
|
|
-
|
|
|
|
|
-def jobrunner(fn,args):
|
|
|
|
|
- print(args)
|
|
|
|
|
- job = fn.delay(*args)
|
|
|
|
|
- print(job.get_id())
|
|
|
|
|
- job_key_list.append(job.get_id())
|
|
|
|
|
- return job
|
|
|
|
|
-
|
|
|
|
|
def utctolocal(date):
|
|
def utctolocal(date):
|
|
|
date=date.replace(tzinfo=tz.tzutc())
|
|
date=date.replace(tzinfo=tz.tzutc())
|
|
|
return date.astimezone(tz.tzlocal())
|
|
return date.astimezone(tz.tzlocal())
|
|
@@ -38,23 +31,32 @@ def result(job_key):
|
|
|
else:
|
|
else:
|
|
|
return ("Process: %s" % json.dumps(job.meta), 202)
|
|
return ("Process: %s" % json.dumps(job.meta), 202)
|
|
|
|
|
|
|
|
-def jobs():
|
|
|
|
|
|
|
+def l_from_ids(ids,status=None):
|
|
|
ret = []
|
|
ret = []
|
|
|
- for j in job_key_list:#q.jobs:
|
|
|
|
|
- try:
|
|
|
|
|
- job = Job.fetch(j, connection=conn)
|
|
|
|
|
- except Exception as e:
|
|
|
|
|
- job_key_list.remove(j)
|
|
|
|
|
- print("%s deleted" % j)
|
|
|
|
|
- continue
|
|
|
|
|
|
|
+ for job_id in ids:
|
|
|
|
|
+ job = Job.fetch(job_id, connection=conn)
|
|
|
out = job.meta
|
|
out = job.meta
|
|
|
- out["id"]=j
|
|
|
|
|
out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
|
|
out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
|
|
|
- out["finished"]=job.is_finished
|
|
|
|
|
|
|
+ out["working"]=True
|
|
|
out["result"]=job.result
|
|
out["result"]=job.result
|
|
|
|
|
+ if status:
|
|
|
|
|
+ out["progress"]=status
|
|
|
ret.append(out)
|
|
ret.append(out)
|
|
|
return ret
|
|
return ret
|
|
|
|
|
|
|
|
|
|
+def jobs():
|
|
|
|
|
+ ret = []
|
|
|
|
|
+ ids = StartedJobRegistry(name='default', connection=conn).get_job_ids()
|
|
|
|
|
+ ret.extend(l_from_ids(ids))
|
|
|
|
|
+ ids = FinishedJobRegistry(name='default', connection=conn).get_job_ids()
|
|
|
|
|
+ ret.extend(l_from_ids(ids,'Finished'))
|
|
|
|
|
+
|
|
|
|
|
+ for q_name in q.all(connection=conn):
|
|
|
|
|
+ ids = q_name.job_ids
|
|
|
|
|
+ ret.extend(l_from_ids(ids,'Waiting'))
|
|
|
|
|
+
|
|
|
|
|
+ return ret
|
|
|
|
|
+
|
|
|
def upload_video(_id,prof,texto1,texto2):
|
|
def upload_video(_id,prof,texto1,texto2):
|
|
|
vid=d.processed(_id)
|
|
vid=d.processed(_id)
|
|
|
ul_file = vid["out_fname"]
|
|
ul_file = vid["out_fname"]
|
|
@@ -65,8 +67,6 @@ def upload_video(_id,prof,texto1,texto2):
|
|
|
#job = jobrunner(upload (titulo, thumb_file, ul_file))
|
|
#job = jobrunner(upload (titulo, thumb_file, ul_file))
|
|
|
#FIXME
|
|
#FIXME
|
|
|
job = upload.delay(titulo, thumb_file, ul_file)
|
|
job = upload.delay(titulo, thumb_file, ul_file)
|
|
|
- job_key_list.append(job.get_id())
|
|
|
|
|
- print(job.get_id())
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
@job('default',connection=conn, timeout=3600*2, result_ttl=3600*2)
|
|
@job('default',connection=conn, timeout=3600*2, result_ttl=3600*2)
|
|
@@ -79,13 +79,18 @@ def upload(title,thumb_path,video_path):
|
|
|
d.insert_uploaded({'title': title,'yt': v_id, 'video':video_path, 'enqueued_at': job_start, 'finished_at': job_end })
|
|
d.insert_uploaded({'title': title,'yt': v_id, 'video':video_path, 'enqueued_at': job_start, 'finished_at': job_end })
|
|
|
return v_id
|
|
return v_id
|
|
|
|
|
|
|
|
-@job('default',connection=conn, timeout=30)
|
|
|
|
|
|
|
+@job('default',connection=conn, timeout=30, result_ttl=135)
|
|
|
def long_job(arg):
|
|
def long_job(arg):
|
|
|
|
|
+ job = get_current_job()
|
|
|
|
|
+ job.meta['type'] = "Test"
|
|
|
|
|
+ job.meta['target'] = 20
|
|
|
|
|
+ job.save()
|
|
|
for i in range(20):
|
|
for i in range(20):
|
|
|
time.sleep(1)
|
|
time.sleep(1)
|
|
|
job = get_current_job()
|
|
job = get_current_job()
|
|
|
job.meta['progress'] = i
|
|
job.meta['progress'] = i
|
|
|
job.save()
|
|
job.save()
|
|
|
|
|
+ print("working")
|
|
|
#push to db, return id if ok. return error if not ok
|
|
#push to db, return id if ok. return error if not ok
|
|
|
return "37"
|
|
return "37"
|
|
|
|
|
|
|
@@ -119,7 +124,7 @@ def crop_video(in_fname,start_time,end_time,out_fname):
|
|
|
|
|
|
|
|
out = run_process(process)
|
|
out = run_process(process)
|
|
|
if out == "":
|
|
if out == "":
|
|
|
- out = { "status": "Proceso finalizado" }
|
|
|
|
|
|
|
+ out = "Proceso finalizado"
|
|
|
|
|
|
|
|
|
|
|
|
|
job_end = datetime.datetime.now()
|
|
job_end = datetime.datetime.now()
|
|
@@ -146,5 +151,7 @@ def run_process(process = None):
|
|
|
job.meta['progress'] = line.rstrip()
|
|
job.meta['progress'] = line.rstrip()
|
|
|
job.save()
|
|
job.save()
|
|
|
else:
|
|
else:
|
|
|
|
|
+ job.meta['progress'] = "Finalizado"
|
|
|
|
|
+ job.save()
|
|
|
break
|
|
break
|
|
|
return ""
|
|
return ""
|