| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164 |
- import os.path
- import datetime
- import json
- import thumb
- import time
- 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 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"]
- titulo = ("%s - %s - %s" % (prof,texto1,texto2))
- thumb_file = thumb.create_thumb(prof,texto1,texto2)
- job = upload.delay(titulo, 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,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,'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,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',
- '-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 = "Proceso finalizado"
- job_end = datetime.datetime.now()
- 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
- })
- 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"
- }
- 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:
- job.meta['progress'] = "Finalizado"
- job.save()
- break
- return ""
|