tasks.py 5.1 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176
  1. from os import stat
  2. import os.path
  3. import datetime
  4. import json
  5. import thumb
  6. import time
  7. from constants import CUT_OUTPUT_PATH
  8. from subprocess import Popen,PIPE
  9. from rq import get_current_job,Queue
  10. from rq.registry import StartedJobRegistry,FinishedJobRegistry
  11. from rq.job import Job
  12. from rq.decorators import job
  13. from worker import conn,q
  14. from dateutil import tz
  15. from db import db
  16. d = db()
  17. date_fmt='%d/%m/%Y, %H:%M'
  18. def utctolocal(date):
  19. date=date.replace(tzinfo=tz.tzutc())
  20. return date.astimezone(tz.tzlocal())
  21. def result(job_key):
  22. """Get job result from key. throws exception if non-existing key"""
  23. job = Job.fetch(job_key, connection=conn)
  24. eq_at=utctolocal(job.enqueued_at)
  25. if job.is_finished:
  26. return ("Result: %s, enqueued at %s" %
  27. (str(job.result),eq_at.strftime(date_fmt)), 200)
  28. else:
  29. return ("Process: %s" % json.dumps(job.meta), 202)
  30. def l_from_ids(ids,status=None):
  31. ret = []
  32. for job_id in ids:
  33. job = Job.fetch(job_id, connection=conn)
  34. out = job.meta
  35. out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
  36. out["working"]=True
  37. #out["queue"]=job.queue
  38. out["result"]=job.result
  39. if status:
  40. out["progress"]=status
  41. ret.append(out)
  42. return ret
  43. def jobs():
  44. ret = []
  45. ids = StartedJobRegistry(name='default', connection=conn).get_job_ids()
  46. ret.extend(l_from_ids(ids))
  47. ids = FinishedJobRegistry(name='default', connection=conn).get_job_ids()
  48. ret.extend(l_from_ids(ids,'Finished'))
  49. for q_name in q.all(connection=conn):
  50. ids = q_name.job_ids
  51. ret.extend(l_from_ids(ids,'Waiting'))
  52. return ret
  53. def upload_video(_id,prof,texto1,texto2):
  54. vid=d.processed(_id)
  55. ul_file = vid["out_fname"]
  56. titulo = ("%s - %s - %s" % (prof,texto1,texto2))
  57. thumb_file = thumb.create_thumb(prof,texto1,texto2)
  58. job = upload.delay(titulo, thumb_file, ul_file)
  59. job.meta["type"]="Youtube"
  60. job.meta["target"]=titulo
  61. @job('default',connection=conn, timeout=3600*2, result_ttl=360)
  62. def upload(title,thumb_path,video_path):
  63. import youtube_upload.main
  64. j=get_current_job()
  65. job_start=datetime.datetime.now()
  66. v_id=youtube_upload.main.main(title,thumb_path,video_path,job=j)
  67. job_end=datetime.datetime.now()
  68. d.insert_uploaded({'title': title,'yt': v_id, 'video':video_path, 'enqueued_at': job_start, 'finished_at': job_end })
  69. j.meta["working"]=False
  70. j.save()
  71. return v_id
  72. @job('default',connection=conn, timeout=30, result_ttl=135)
  73. def long_job(arg):
  74. job = get_current_job()
  75. job.meta['type'] = "Test"
  76. job.meta['target'] = 20
  77. job.save()
  78. for i in range(20):
  79. time.sleep(1)
  80. job = get_current_job()
  81. job.meta['progress'] = i
  82. job.save()
  83. print("working")
  84. #push to db, return id if ok. return error if not ok
  85. return "37"
  86. @job('default',connection=conn, timeout=360*2, result_ttl=360)
  87. def crop_video(in_fname,start_time,end_time,out_fname):
  88. job = get_current_job()
  89. job.meta["type"]="Crop"
  90. job.save()
  91. job_start = datetime.datetime.now()
  92. if not os.path.isfile(in_fname):
  93. return { "status":"error",
  94. "error": "El archivo %s no existe" % in_fname
  95. }
  96. out_fname=os.path.join(CUT_OUTPUT_PATH,out_fname)
  97. if os.path.isfile(out_fname):
  98. return { "status":"error",
  99. "error": "El archivo %s existe" % out_fname
  100. }
  101. process = [ 'ffmpeg',
  102. '-loglevel', 'error',
  103. '-stats',
  104. '-analyzeduration', '16M',
  105. '-i', in_fname,
  106. '-ss', str(start_time),
  107. '-t', str(end_time),
  108. '-c', 'copy',
  109. '-movflags', '+faststart',
  110. '-f','mp4',
  111. out_fname
  112. ]
  113. out = run_process(process)
  114. if out == "":
  115. out = "Proceso finalizado"
  116. #
  117. test_broken = [ "ffprobe", "-loglevel", "warning", out_fname ]
  118. is_broken = Popen(test_broken, stderr=PIPE)
  119. l = is_broken.stderr.readline()
  120. if len(l) > 3:
  121. print(l)
  122. job_end = datetime.datetime.now()
  123. filesize = stat(out_fname).st_size
  124. filesize = int(filesize/(1024*1024))
  125. d.insert_processed({
  126. 'job_start': job_start,
  127. 'job_end': job_end,
  128. 'start_time':start_time,
  129. 'end_time': end_time,
  130. 'in_fname': in_fname,
  131. 'out_fname': out_fname,
  132. 'filesize': filesize
  133. })
  134. job.meta["working"]=False
  135. job.save()
  136. return out
  137. def run_process(process = None):
  138. if process is None or type(process) is not list:
  139. return { "status": "error",
  140. "error": "Argumento process invalido"
  141. }
  142. proc=Popen(process, stderr=PIPE, universal_newlines=True)
  143. job = get_current_job()
  144. while True:
  145. line = proc.stderr.readline()
  146. if line != '' and line !=b'':
  147. job.meta['progress'] = line.rstrip()
  148. job.save()
  149. else:
  150. job.meta['progress'] = "Finalizado"
  151. job.save()
  152. break
  153. return ""