tasks.py 4.3 KB

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