unpacker.py 4.2 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116
  1. import os
  2. import json
  3. import socket
  4. import sys
  5. import time
  6. import tempfile
  7. from threading import Thread
  8. from subprocess import check_output
  9. from bookworm import s3
  10. from bookworm.logger import log, setup_logger
  11. from bookworm.constants import RAW_FILE_BUCKET, PROCESSED_FILE_BUCKET, UNPACKABLE_EXTENSIONS, REDIS_UNPACK_FILE, REDIS_STATE_KEY, REDIS_STEP_KEY
  12. import redis
  13. def should_unpack(fname):
  14. fname = fname.lower()
  15. return fname.endswith('rar') or fname.endswith('zip')
  16. def archive_contents(fd):
  17. to_extract = {}
  18. contents = check_output(['lsar', '-j', fd.name]).decode('utf-8')
  19. contents = json.loads(contents)
  20. log.debug(contents['lsarContents'])
  21. for f in contents['lsarContents']:
  22. fname = f['XADFileName']
  23. if any([extension in fname.lower() for extension in UNPACKABLE_EXTENSIONS]):
  24. log.info('Extracting %s from the archive', fname)
  25. to_extract[f['XADIndex']] = fname
  26. return to_extract
  27. def convert_to_mobi(orig_fname) -> bytes:
  28. with tempfile.NamedTemporaryFile(suffix='.mobi') as fd:
  29. log.info('Converting %s to mobi at %s', orig_fname, fd.name)
  30. check_output(['ebook-convert', orig_fname, fd.name, '--output-profile=kindle_pw'])
  31. log.info('Done converting')
  32. buff = open(fd.name, 'rb').read()
  33. return buff
  34. def store_file(s3client, fname, file_contents):
  35. log.info('Puttin in s3 under bucket %s with key %s', PROCESSED_FILE_BUCKET, fname)
  36. s3client.put_object(Body=file_contents, Bucket=PROCESSED_FILE_BUCKET, Key=fname)
  37. log.info('Put in s3 under bucket %s with key %s', PROCESSED_FILE_BUCKET, fname)
  38. def delete_raw_file(s3client, s3key):
  39. log.info('Deleting %s from %s', s3key, RAW_FILE_BUCKET)
  40. s3client.delete_object(Bucket=RAW_FILE_BUCKET, Key=s3key)
  41. def unpack_and_convert(job_key, s3key, s3client, redis):
  42. redis.hset(job_key, REDIS_STEP_KEY, 'UNPACKING')
  43. unpacked_files = unpack(s3key, s3client, redis)
  44. redis.hset(job_key, REDIS_STEP_KEY, 'UNPACK_DONE')
  45. log.info('Done unpacking job %s', job_key)
  46. for fname, data in unpacked_files:
  47. if not fname.lower().endswith('mobi'):
  48. redis.hset(job_key, REDIS_STEP_KEY, 'CONVERTING')
  49. redis.hset(job_key, REDIS_STATE_KEY, fname)
  50. log.info('Asked to store %s, need to convert first', fname)
  51. fname_no_ext, ext = os.path.splitext(fname)
  52. with tempfile.NamedTemporaryFile(suffix=ext) as original:
  53. original.write(data)
  54. original.flush()
  55. converted_data = convert_to_mobi(original.name)
  56. data = converted_data
  57. fname = fname_no_ext + '.mobi'
  58. store_file(s3client, fname, data)
  59. delete_raw_file(s3client, s3key)
  60. redis.delete(job_key)
  61. def unpack(s3key, s3client, redis):
  62. log.info('Got a request to unpack %s', s3key)
  63. data = s3client.get_object(Key=s3key, Bucket=RAW_FILE_BUCKET)
  64. with tempfile.NamedTemporaryFile() as fd:
  65. raw_file_contents = data['Body'].read()
  66. fd.write(raw_file_contents)
  67. fd.flush()
  68. if not should_unpack(s3key):
  69. log.info("Not unpacking %s", s3key)
  70. return [(s3key, raw_file_contents)]
  71. to_extract = archive_contents(fd)
  72. if not to_extract:
  73. log.info(contents['lsarContents'])
  74. log.error("Could not find any valid file")
  75. return []
  76. ret = []
  77. for index, fname in to_extract.items():
  78. log.info('Processing %s %s', index, fname)
  79. file_contents = check_output(['unar', '-o', '-', '-i', fd.name, str(index)])
  80. log.info('Got %d bytes', len(file_contents))
  81. ret.append((fname, file_contents))
  82. return ret
  83. def main():
  84. r = redis.StrictRedis(host='localhost', port=6379)
  85. setup_logger()
  86. s3client = s3.client()
  87. while True:
  88. log.info('Waiting for message...')
  89. topic, message = r.blpop(REDIS_UNPACK_FILE)
  90. log.info('got message: %s', message)
  91. params = json.loads(message.decode('utf-8'))
  92. log.info('params for unpacker: %s', params)
  93. params['s3client'] = s3client
  94. params['redis'] = r
  95. t = Thread(target=unpack_and_convert, kwargs=params)
  96. t.daemon = True
  97. t.start()
  98. main()