tasks.py 4.7 KB

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