tasks.py 4.7 KB

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