import os.path import time import json from constants import CUT_OUTPUT_PATH from subprocess import Popen,PIPE from rq import get_current_job,Queue from rq.job import Job from rq.decorators import job 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): for i in range(20): time.sleep(1) job = get_current_job() job.meta['progress'] = i job.save() #push to db, return id if ok. return error if not ok return "37" @job('default',connection=conn, timeout=3600*2, result_ttl=3600*2) def upload_video(title,thumb_path,video_path): import youtube_upload.main j=get_current_job() return youtube_upload.main.main(title,thumb_path,video_path,job=j) @job('default',connection=conn, timeout=360*2, result_ttl=360*2) def cut_video(in_fname,start_time,end_time,out_fname): if not os.path.isfile(in_fname): return { "status":"error", "error": "El archivo %s no existe" % in_fname } out_fname=os.path.join(CUT_OUTPUT_PATH,out_fname) if os.path.isfile(out_fname): return { "status":"error", "error": "El archivo %s existe" % out_fname } process = [ 'ffmpeg', '-loglevel', 'error', '-stats', '-i', in_fname, '-ss', str(start_time), '-t', str(end_time), '-c', 'copy', '-movflags', '+faststart', '-f','mp4', out_fname ] out = run_process(process) if out == "": out = { "status": "Proceso finalizado" } return out def run_process(process = None): if process is None or type(process) is not list: return { "status": "error", "error": "Argumento process invalido" } proc=Popen(process, stderr=PIPE, universal_newlines=True) while True: line = proc.stderr.readline() if line != '' and line !=b'': job = get_current_job() job.meta['progress'] = line.rstrip() job.save() else: break return ""