from os import stat import os.path import datetime import json import thumb import time import requests from constants import CUT_OUTPUT_PATH from subprocess import Popen,PIPE from rq import get_current_job,Queue from rq.registry import StartedJobRegistry,FinishedJobRegistry from rq.job import Job from rq.decorators import job from worker import conn,q from dateutil import tz from db import db d = db() date_fmt='%d/%m/%Y, %H:%M' 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 cargar_plataforma(video_id): data = d.videoData(video_id) prof = None if "profesor" in data: prof = data["profesor"] tsplit = data["title"].split("-") texto1 = tsplit[1].strip() texto2 = tsplit[2].strip() tmp = { "video": data["yt"], "diapos": "[]", "titulo": "%s - %s" % (texto1, texto2), "auth_token": "4b494d6c7ffc74d387d8cfdef8ca1691" } if prof is not None: tmp["profesor"] = prof r = requests.post('https://plataforma.especificosba.com.ar/back/crearVideo.php', json = tmp ); if r.status_code == 200: return {"status":"Ok"} else: return {"status": r.text, "postdata": tmp, "ret": r.status_code} def l_from_ids(ids,status=None): ret = [] for job_id in ids: job = Job.fetch(job_id, connection=conn) out = job.meta out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt) out["working"]=True #out["queue"]=job.queue out["result"]=job.result if status: out["progress"]=status ret.append(out) 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): vid=d.processed(_id) ul_file = vid["out_fname"] profname = prof try: profname = ("%s %s" % (prof["nombre"], prof["apellido"])) except: pass titulo = ("%s - %s - %s" % (profname,texto1,texto2)) thumb_file = thumb.create_thumb(profname,texto1,texto2) job = upload.delay(titulo, prof, thumb_file, ul_file) job.meta["type"]="Youtube" job.meta["target"]=titulo @job('default',connection=conn, timeout=3600*2, result_ttl=360) def upload(title,prof,thumb_path,video_path): import youtube_upload.main j=get_current_job() job_start=datetime.datetime.now() v_id=youtube_upload.main.main(title,thumb_path,video_path,job=j) job_end=datetime.datetime.now() d.insert_uploaded({'title': title,'profesor':prof,'yt': v_id, 'video':video_path, 'enqueued_at': job_start, 'finished_at': job_end }) j.meta["working"]=False j.save() return v_id @job('default',connection=conn, timeout=30, result_ttl=135) def long_job(arg): job = get_current_job() job.meta['type'] = "Test" job.meta['target'] = 20 job.save() for i in range(20): time.sleep(1) job = get_current_job() job.meta['progress'] = i job.save() print("working") #push to db, return id if ok. return error if not ok return "37" @job('default',connection=conn, timeout=360*2, result_ttl=360) def crop_video(in_fname,start_time,end_time,censuras,out_fname): job = get_current_job() job.meta["type"]="Crop" job.save() job_start = datetime.datetime.now() 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', '-analyzeduration', '8M', '-i', in_fname ] if censuras is not None and len(censuras) > 0: s = "" print(censuras) for c in censuras: s += "volume=enable='between(t,%d,%d)':volume=0, " % (int(c["inicio"]),int(c["fin"])) s=s[:-2] #borro ", " del final process = process + [ "-c:a", "libfdk_aac", "-af", s ] else: process += ["-c:a", "copy"] process += [ '-ss', str(start_time), '-to', str(end_time), '-c:v', 'copy', '-movflags', '+faststart', '-f', 'mp4', out_fname ] out = run_process(process) if out == "": out = "Proceso finalizado" test_broken = [ "ffprobe", "-loglevel", "warning", out_fname ] is_broken = Popen(test_broken, stderr=PIPE) l = is_broken.stderr.readline() if len(l) > 3: print(l) job_end = datetime.datetime.now() filesize = stat(out_fname).st_size filesize = int(filesize/(1024*1024)) d.insert_processed({ 'job_start': job_start, 'job_end': job_end, 'start_time':start_time, 'end_time': end_time, 'in_fname': in_fname, 'out_fname': out_fname, 'filesize': filesize }) job.meta["working"]=False job.save() return out def run_process(process = None): if process is None or type(process) is not list: return { "status": "error", "error": "Argumento process invalido" } print(" ".join(process)) proc=Popen(process, stderr=PIPE, universal_newlines=True) job = get_current_job() while True: line = proc.stderr.readline() if line != '' and line !=b'': job.meta['progress'] = line.rstrip() job.save() else: job.meta['progress'] = "Finalizado" job.save() break return ""