tasks.py 6.3 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217
  1. from os import stat
  2. import os.path
  3. import datetime
  4. import json
  5. import thumb
  6. import time
  7. import requests
  8. from constants import CUT_OUTPUT_PATH
  9. from subprocess import Popen,PIPE
  10. from rq import get_current_job,Queue
  11. from rq.registry import StartedJobRegistry,FinishedJobRegistry
  12. from rq.job import Job
  13. from rq.decorators import job
  14. from worker import conn,q
  15. from dateutil import tz
  16. from db import db
  17. d = db()
  18. date_fmt='%d/%m/%Y, %H:%M'
  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 cargar_plataforma(video_id):
  32. data = d.videoData(video_id)
  33. prof = None
  34. if "profesor" in data:
  35. prof = data["profesor"]
  36. tsplit = data["title"].split("-")
  37. texto1 = tsplit[1].strip()
  38. texto2 = tsplit[2].strip()
  39. tmp = { "video": data["yt"],
  40. "diapos": "[]",
  41. "titulo": "%s - %s" % (texto1, texto2),
  42. "auth_token": "4b494d6c7ffc74d387d8cfdef8ca1691"
  43. }
  44. if prof is not None:
  45. tmp["profesor"] = prof
  46. r = requests.post('https://plataforma.especificosba.com.ar/back/crearVideo.php', json = tmp );
  47. if r.status_code == 200:
  48. return {"status":"Ok"}
  49. else:
  50. return {"status": r.text, "postdata": tmp, "ret": r.status_code}
  51. def l_from_ids(ids,status=None):
  52. ret = []
  53. for job_id in ids:
  54. job = Job.fetch(job_id, connection=conn)
  55. out = job.meta
  56. out["eq"]=utctolocal(job.enqueued_at).strftime(date_fmt)
  57. out["working"]=True
  58. #out["queue"]=job.queue
  59. out["result"]=job.result
  60. if status:
  61. out["progress"]=status
  62. ret.append(out)
  63. return ret
  64. def jobs():
  65. ret = []
  66. ids = StartedJobRegistry(name='default', connection=conn).get_job_ids()
  67. ret.extend(l_from_ids(ids))
  68. ids = FinishedJobRegistry(name='default', connection=conn).get_job_ids()
  69. ret.extend(l_from_ids(ids,'Finished'))
  70. for q_name in q.all(connection=conn):
  71. ids = q_name.job_ids
  72. ret.extend(l_from_ids(ids,'Waiting'))
  73. return ret
  74. def upload_video(_id,prof,texto1,texto2):
  75. vid=d.processed(_id)
  76. ul_file = vid["out_fname"]
  77. profname = prof
  78. try:
  79. profname = ("%s %s" % (prof["nombre"], prof["apellido"]))
  80. except:
  81. pass
  82. titulo = ("%s - %s - %s" % (profname,texto1,texto2))
  83. thumb_file = thumb.create_thumb(profname,texto1,texto2)
  84. job = upload.delay(titulo, prof, thumb_file, ul_file)
  85. job.meta["type"]="Youtube"
  86. job.meta["target"]=titulo
  87. @job('default',connection=conn, timeout=3600*2, result_ttl=360)
  88. def upload(title,prof,thumb_path,video_path):
  89. import youtube_upload.main
  90. j=get_current_job()
  91. job_start=datetime.datetime.now()
  92. v_id=youtube_upload.main.main(title,thumb_path,video_path,job=j)
  93. job_end=datetime.datetime.now()
  94. d.insert_uploaded({'title': title,'profesor':prof,'yt': v_id, 'video':video_path, 'enqueued_at': job_start, 'finished_at': job_end })
  95. j.meta["working"]=False
  96. j.save()
  97. return v_id
  98. @job('default',connection=conn, timeout=30, result_ttl=135)
  99. def long_job(arg):
  100. job = get_current_job()
  101. job.meta['type'] = "Test"
  102. job.meta['target'] = 20
  103. job.save()
  104. for i in range(20):
  105. time.sleep(1)
  106. job = get_current_job()
  107. job.meta['progress'] = i
  108. job.save()
  109. print("working")
  110. #push to db, return id if ok. return error if not ok
  111. return "37"
  112. @job('default',connection=conn, timeout=360*2, result_ttl=360)
  113. def crop_video(in_fname,start_time,end_time,censuras,out_fname):
  114. job = get_current_job()
  115. job.meta["type"]="Crop"
  116. job.save()
  117. job_start = datetime.datetime.now()
  118. if not os.path.isfile(in_fname):
  119. return { "status":"error",
  120. "error": "El archivo %s no existe" % in_fname
  121. }
  122. out_fname=os.path.join(CUT_OUTPUT_PATH,out_fname)
  123. if os.path.isfile(out_fname):
  124. return { "status":"error",
  125. "error": "El archivo %s existe" % out_fname
  126. }
  127. process = [ 'ffmpeg',
  128. '-analyzeduration', '8M',
  129. '-i', in_fname
  130. ]
  131. if censuras is not None and len(censuras) > 0:
  132. s = ""
  133. print(censuras)
  134. for c in censuras:
  135. s += "volume=enable='between(t,%d,%d)':volume=0, " % (int(c["inicio"]),int(c["fin"]))
  136. s=s[:-2] #borro ", " del final
  137. process = process + [ "-c:a", "libfdk_aac", "-af", s ]
  138. else:
  139. process += ["-c:a", "copy"]
  140. process += [
  141. '-ss', str(start_time),
  142. '-to', str(end_time),
  143. '-c:v', 'copy',
  144. '-movflags', '+faststart',
  145. '-f', 'mp4', out_fname
  146. ]
  147. out = run_process(process)
  148. if out == "":
  149. out = "Proceso finalizado"
  150. test_broken = [ "ffprobe", "-loglevel", "warning", out_fname ]
  151. is_broken = Popen(test_broken, stderr=PIPE)
  152. l = is_broken.stderr.readline()
  153. if len(l) > 3:
  154. print(l)
  155. job_end = datetime.datetime.now()
  156. filesize = stat(out_fname).st_size
  157. filesize = int(filesize/(1024*1024))
  158. d.insert_processed({
  159. 'job_start': job_start,
  160. 'job_end': job_end,
  161. 'start_time':start_time,
  162. 'end_time': end_time,
  163. 'in_fname': in_fname,
  164. 'out_fname': out_fname,
  165. 'filesize': filesize
  166. })
  167. job.meta["working"]=False
  168. job.save()
  169. return out
  170. def run_process(process = None):
  171. if process is None or type(process) is not list:
  172. return { "status": "error",
  173. "error": "Argumento process invalido"
  174. }
  175. print(" ".join(process))
  176. proc=Popen(process, stderr=PIPE, universal_newlines=True)
  177. job = get_current_job()
  178. while True:
  179. line = proc.stderr.readline()
  180. if line != '' and line !=b'':
  181. job.meta['progress'] = line.rstrip()
  182. job.save()
  183. else:
  184. job.meta['progress'] = "Finalizado"
  185. job.save()
  186. break
  187. return ""