tasks.py 3.6 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127
  1. import os.path
  2. import datetime
  3. import json
  4. from constants import CUT_OUTPUT_PATH
  5. from subprocess import Popen,PIPE
  6. from rq import get_current_job,Queue
  7. from rq.job import Job
  8. from rq.decorators import job
  9. from worker import conn,q
  10. from dateutil import tz
  11. from db import db
  12. d = db()
  13. job_key_list = []
  14. date_fmt='%d/%m/%Y, %H:%M'
  15. def jobrunner(fn,args):
  16. job = fn.delay(args)
  17. job_key_list.append(job.get_id())
  18. return job
  19. def utctolocal(date):
  20. date=date.replace(tzinfo=tz.tzutc())
  21. return date.astimezone(tz.tzlocal())
  22. def result(job_key):
  23. """Get job result from key. throws exception if non-existing key"""
  24. job = Job.fetch(job_key, connection=conn)
  25. eq_at=utctolocal(job.enqueued_at)
  26. if job.is_finished:
  27. return ("Result: %s, enqueued at %s" %
  28. (str(job.result),eq_at.strftime(date_fmt)), 200)
  29. else:
  30. return ("Process: %s" % json.dumps(job.meta), 202)
  31. def jobs():
  32. ret = []
  33. for j in job_key_list:#q.jobs:
  34. try:
  35. job = Job.fetch(j, connection=conn)
  36. except Exception as e:
  37. job_key_list.remove(j)
  38. print("%s deleted" % j)
  39. out = job.meta
  40. out["id"]=j
  41. out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
  42. out["finished"]=job.is_finished
  43. out["result"]=job.result
  44. ret.append(out)
  45. return ret
  46. @job('default',connection=conn, timeout=30)
  47. def long_job(arg):
  48. for i in range(20):
  49. time.sleep(1)
  50. job = get_current_job()
  51. job.meta['progress'] = i
  52. job.save()
  53. #push to db, return id if ok. return error if not ok
  54. return "37"
  55. @job('default',connection=conn, timeout=3600*2, result_ttl=3600*2)
  56. def upload_video(title,thumb_path,video_path):
  57. import youtube_upload.main
  58. j=get_current_job()
  59. return youtube_upload.main.main(title,thumb_path,video_path,job=j)
  60. @job('default',connection=conn, timeout=360*2, result_ttl=360*2)
  61. def cut_video(in_fname,start_time,end_time,out_fname):
  62. job_start = datetime.datetime.now()
  63. if not os.path.isfile(in_fname):
  64. return { "status":"error",
  65. "error": "El archivo %s no existe" % in_fname
  66. }
  67. out_fname=os.path.join(CUT_OUTPUT_PATH,out_fname)
  68. if os.path.isfile(out_fname):
  69. return { "status":"error",
  70. "error": "El archivo %s existe" % out_fname
  71. }
  72. process = [ 'ffmpeg',
  73. '-loglevel', 'error',
  74. '-stats',
  75. '-i', in_fname,
  76. '-ss', str(start_time),
  77. '-t', str(end_time),
  78. '-c', 'copy',
  79. '-movflags', '+faststart',
  80. '-f','mp4',
  81. out_fname
  82. ]
  83. out = run_process(process)
  84. if out == "":
  85. out = { "status": "Proceso finalizado" }
  86. job_end = datetime.datetime.now()
  87. d.insert_processed({
  88. 'job_start': job_start,
  89. 'job_end': job_end,
  90. 'start_time':start_time,
  91. 'end_time': end_time,
  92. 'in_fname': in_fname,
  93. 'out_fname': out_fname
  94. })
  95. return out
  96. def run_process(process = None):
  97. if process is None or type(process) is not list:
  98. return { "status": "error",
  99. "error": "Argumento process invalido"
  100. }
  101. proc=Popen(process, stderr=PIPE, universal_newlines=True)
  102. while True:
  103. line = proc.stderr.readline()
  104. if line != '' and line !=b'':
  105. job = get_current_job()
  106. job.meta['progress'] = line.rstrip()
  107. job.save()
  108. else:
  109. break
  110. return ""