tasks.py 3.1 KB

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