David 10 anni fa
parent
commit
7e8d98b1c0
4 ha cambiato i file con 73 aggiunte e 22 eliminazioni
  1. 4 0
      back/db.py
  2. 41 4
      back/tasks.py
  3. 25 16
      back/web.py
  4. 3 2
      back/worker.py

+ 4 - 0
back/db.py

@@ -15,7 +15,11 @@ class db():
         self.raw_files = self.client[self.DB]["raw"]
         self.cut_files = self.client[self.DB]["cut"]
         self.ul_files  = self.client[self.DB]["ul"]
+        self.jobs      = self.client[self.DB]["jobs"]
 
+    def getJobs(self):
+        ret = [ jsonable(r) for r in self.jobs.find() ]
+        return ret
 
     def raw(self, _id=None):
         ret = None

+ 41 - 4
back/tasks.py

@@ -1,11 +1,48 @@
 import os.path
 import time
+import json
 
 from constants import CUT_OUTPUT_PATH
 from subprocess import Popen,PIPE
-from rq import get_current_job
+from rq import get_current_job,Queue
+from rq.job import Job
 from rq.decorators import job
-from worker import conn
+from worker import conn,q
+from dateutil import tz
+
+job_key_list = []
+date_fmt='%d/%m/%Y, %H:%M'
+
+
+def jobrunner(fn,args):
+    job = fn.delay(args)
+    job_key_list.append(job.get_id())
+    return job
+
+def utctolocal(date):
+    date=date.replace(tzinfo=tz.tzutc())
+    return date.astimezone(tz.tzlocal())
+
+def result(job_key):
+    """Get job result from key. throws exception if non-existing key"""
+    job = Job.fetch(job_key, connection=conn)
+    eq_at=utctolocal(job.enqueued_at)
+    if job.is_finished:
+        return ("Result: %s, enqueued at %s" %
+                    (str(job.result),eq_at.strftime(date_fmt)), 200)
+    else:
+        return ("Process: %s" % json.dumps(job.meta), 202)
+
+def jobs():
+    ret = []
+    for j in job_key_list:#q.jobs:
+        job = Job.fetch(j, connection=conn)
+        out = job.meta
+        out["id"]=j
+        out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
+        out["finished"]=job.is_finished
+        ret.append(out)
+    return ret
 
 @job('default',connection=conn, timeout=30)
 def long_job(arg):
@@ -41,8 +78,8 @@ def cut_video(in_fname,start_time,end_time,out_fname):
                 '-loglevel', 'error',
                 '-stats',
                 '-i', in_fname,
-                '-ss', start_time,
-                '-t', end_time,
+                '-ss', str(start_time),
+                '-t', str(end_time),
                 '-c', 'copy',
                 '-movflags', '+faststart',
                 '-f','mp4',

+ 25 - 16
back/web.py

@@ -7,10 +7,10 @@ import tasks
 from db import db
 from flask import Flask, jsonify, request
 from rq import Queue
-from rq.job import Job
-from worker import conn
+
 app = Flask(__name__)
 d=db()
+
 @app.route(BASE_PATH + '/')
 def index():
     return "Hello world"
@@ -23,22 +23,37 @@ def raw_list():
 @app.route(BASE_PATH + '/raw/<video_id>')
 def raw_file(video_id):
     ret=d.raw(video_id)
-    ret["file"]=ret["file"].split("front")[1] #FIXME
+    ret["file"]=ret["file"].split("front")[1] #FIXME, path relativo al webserver
     return jsonify(ret)
 
 @app.route(BASE_PATH + '/cut/', methods=["POST"])
 def cut():
     ret={"status":"Procesando"}
+    today=datetime.date.today()
     data = request.get_json()
     video=d.raw(data["id"])
-    out_filename="%s_%s_%s.mp4" % (data["curso"], datetime.date.today().strftime('%Y%m%d'), data["desc"])
-    job = tasks.cut_video.delay(video["file"],data["inicio"],data["fin"],out_filename)
+
+    out_filename="%s_%s_%s.mp4" % (data["curso"],
+                                today.strftime('%Y%m%d'),
+                                data["desc"])
+
+    job = tasks.cut_video.delay(video["file"],
+                                str(data["inicio"]),
+                                str(data["fin"]),
+                                out_filename)
     ret["job_id"]=job.get_id()
     return jsonify(ret)
 
+@app.route(BASE_PATH + '/jobs/')
+def jobs():
+    working = tasks.jobs()
+    finished = d.getJobs()
+    return json.dumps({"working":working, "finished":finished})
+
 @app.route(BASE_PATH + '/job')
 def job():
-#    job = tasks.long_job.delay("test")
+    job = tasks.jobrunner(tasks.long_job, ("test",))
+    #job = tasks.long_job.delay("test")
 #    job = tasks.upload_video.delay("titulo","/home/david/photo_2016-07-10_13-49-56.jpg","/home/david/movie-1467033043.flv")
 
     ret = {"id":job.get_id()}
@@ -48,18 +63,12 @@ def job():
 
 @app.route(BASE_PATH+"/results/<job_key>", methods=['GET'])
 def get_results(job_key):
-
-    print(job_key)
     try:
-        job = Job.fetch(job_key, connection=conn)
-    except:
+        ret=tasks.result(job_key)
+    except Exception as e:
+        print(e)
         return "Error. Probablemente job_key es invalido", 400
-
-
-    if job.is_finished:
-        return "Result: %s" % str(job.result), 200
-    else:
-        return "Process: %s" % json.dumps(job.meta) , 202
+    return ret
 
 
 @app.route(BASE_PATH + '/uploaded')

+ 3 - 2
back/worker.py

@@ -6,12 +6,13 @@ from rq import Worker, Queue, Connection
 
 listen = ['default']
 
-redis_url = os.getenv('REDISTOGO_URL', 'redis://localhost:6379')
+redis_url = 'redis://localhost:6379'
 
 conn = redis.from_url(redis_url)
+q = Queue('default', connection=conn)
 
 if __name__ == '__main__':
     with Connection(conn):
-        worker = Worker(list(map(Queue, listen)))
+        worker = Worker([q])
         worker.work()