tasks.py 5.7 KB

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