소스 검색

add status transitions via redis

david 7 년 전
부모
커밋
851750a5ed
4개의 변경된 파일30개의 추가작업 그리고 23개의 파일을 삭제
  1. 2 0
      bookworm/constants.py
  2. 6 5
      bookworm/file_fetcher.py
  3. 3 3
      bookworm/ircclient.py
  4. 19 15
      bookworm/unpacker.py

+ 2 - 0
bookworm/constants.py

@@ -6,3 +6,5 @@ UNPACKABLE_EXTENSIONS = ['epub', 'mobi', 'azw3']
 REDIS_BOOK_COMMANDS = 'BOOK_COMMANDS'
 REDIS_FETCH_FILE = 'FETCH_FILE'
 REDIS_UNPACK_FILE = 'UNPACK_FILE'
+REDIS_STEP_KEY = 'STEP'
+REDIS_STATE_KEY = 'STATE'

+ 6 - 5
bookworm/file_fetcher.py

@@ -3,7 +3,7 @@ import socket
 import sys
 import time
 from threading import Thread
-from bookworm.constants import RAW_FILE_BUCKET, REDIS_UNPACK_FILE, REDIS_FETCH_FILE
+from bookworm.constants import RAW_FILE_BUCKET, REDIS_UNPACK_FILE, REDIS_FETCH_FILE, REDIS_STATE_KEY, REDIS_STEP_KEY
 from bookworm.logger import log, setup_logger
 from bookworm import s3
 
@@ -16,6 +16,7 @@ def netcat(filename, ip, port, size, job_key, s3client, redis):
     s.connect((ip, port))
     log.info("Receiving file %s", filename)
 
+    redis.hset(job_key, REDIS_STEP_KEY, 'DOWNLOADING')
     buff = b''
     count = 0
     last_perc = 0
@@ -29,17 +30,17 @@ def netcat(filename, ip, port, size, job_key, s3client, redis):
         perc = int(100 * count / size)
         if perc % 5 == 0 and perc != last_perc:
             log.info("Download percentage: %d", perc)
-            # TODO set_job_state(job_key, 'DOWNLOADING', "%d%%" % perc)
+            redis.hset(job_key, REDIS_STATE_KEY, str(perc))
             last_perc = perc
         if count >= size:
             break
     log.info("Download complete")
     s.close()
     log.info("Putting file in s3")
-    # TODO set_job_state(job_key, 'DOWNLOAD_DONE', job_key)
-    s3client.put_object(Body=buff, Bucket=RAW_FILE_BUCKET, Key=job_key)
+    redis.hset(job_key, REDIS_STATE_KEY, '100')
+    s3client.put_object(Body=buff, Bucket=RAW_FILE_BUCKET, Key=filename)
     log.info("File %s in s3 with key %s", filename, job_key)
-    redis.rpush(REDIS_UNPACK_FILE, json.dumps({'job_key': job_key}))
+    redis.rpush(REDIS_UNPACK_FILE, json.dumps({'job_key': job_key, 's3key': filename}))
 
 def main():
     r = redis.StrictRedis(host='localhost', port=6379)

+ 3 - 3
bookworm/ircclient.py

@@ -6,7 +6,7 @@ import shlex
 
 import utils
 from threading import Thread
-from bookworm.constants import IRC_TIME_TO_FIRST_COMMAND, IRC_CHANNEL, REDIS_BOOK_COMMANDS, REDIS_FETCH_FILE
+from bookworm.constants import IRC_TIME_TO_FIRST_COMMAND, IRC_CHANNEL, REDIS_BOOK_COMMANDS, REDIS_FETCH_FILE, REDIS_STEP_KEY
 from bookworm.logger import log, setup_logger
 
 import redis
@@ -53,7 +53,7 @@ class IRCClient(irc.client.SimpleIRCClient):
 
             bot = command['bot'].strip()
             book = command['book'].strip()
-            # set_job_state(job_key, 'waiting', time.time())
+            self.r.hset('book_'+book, REDIS_STEP_KEY, 'REQUESTED')
             self.connection.privmsg(self.target, f'!{bot} {book}')
 
     def on_pubmsg(self, connection, event):
@@ -84,7 +84,7 @@ class IRCClient(irc.client.SimpleIRCClient):
         filename, peer_address, peer_port, size = parts
         peer_address = irc.client.ip_numstr_to_quad(peer_address)
         peer_port = int(peer_port)
-        job_key = filename
+        job_key = 'book_' + filename
         data = json.dumps({'ip': peer_address,
                            'port': peer_port,
                            'size': int(size),

+ 19 - 15
bookworm/unpacker.py

@@ -8,7 +8,7 @@ from threading import Thread
 from subprocess import check_output
 from bookworm import s3
 from bookworm.logger import log, setup_logger
-from bookworm.constants import RAW_FILE_BUCKET, PROCESSED_FILE_BUCKET, UNPACKABLE_EXTENSIONS, REDIS_UNPACK_FILE
+from bookworm.constants import RAW_FILE_BUCKET, PROCESSED_FILE_BUCKET, UNPACKABLE_EXTENSIONS, REDIS_UNPACK_FILE, REDIS_STATE_KEY, REDIS_STEP_KEY
 import redis
 
 def should_unpack(fname):
@@ -41,17 +41,20 @@ def store_file(s3client, fname, file_contents):
     s3client.put_object(Body=file_contents, Bucket=PROCESSED_FILE_BUCKET, Key=fname)
     log.info('Put in s3 under bucket %s with key %s', PROCESSED_FILE_BUCKET, fname)
 
-def delete_raw_file(s3client, job_key):
-    log.info('Deleting %s from %s', job_key, RAW_FILE_BUCKET)
-    s3client.delete_object(Bucket=RAW_FILE_BUCKET, Key=job_key)
+def delete_raw_file(s3client, s3key):
+    log.info('Deleting %s from %s', s3key, RAW_FILE_BUCKET)
+    s3client.delete_object(Bucket=RAW_FILE_BUCKET, Key=s3key)
 
-def unpack_and_convert(job_key, s3client, redis):
-    unpacked_files = unpack(job_key, s3client, redis)
-    # TODO set_job_state(job_key, 'UNPACK_DONE', job_key)
+def unpack_and_convert(job_key, s3key, s3client, redis):
+    redis.hset(job_key, REDIS_STEP_KEY, 'UNPACKING')
+    unpacked_files = unpack(s3key, s3client, redis)
+    redis.hset(job_key, REDIS_STEP_KEY, 'UNPACK_DONE')
     log.info('Done unpacking job %s', job_key)
 
     for fname, data in unpacked_files:
         if not fname.lower().endswith('mobi'):
+            redis.hset(job_key, REDIS_STEP_KEY, 'CONVERTING')
+            redis.hset(job_key, REDIS_STATE_KEY, fname)
             log.info('Asked to store %s, need to convert first', fname)
             fname_no_ext, ext = os.path.splitext(fname)
             with tempfile.NamedTemporaryFile(suffix=ext) as original:
@@ -62,19 +65,20 @@ def unpack_and_convert(job_key, s3client, redis):
             fname = fname_no_ext + '.mobi'
         store_file(s3client, fname, data)
 
-    delete_raw_file(s3client, job_key)
+    delete_raw_file(s3client, s3key)
+    redis.delete(job_key)
 
-def unpack(job_key, s3client, redis):
-    log.info('Got a request to unpack %s', job_key)
-    data = s3client.get_object(Key=job_key, Bucket=RAW_FILE_BUCKET)
+def unpack(s3key, s3client, redis):
+    log.info('Got a request to unpack %s', s3key)
+    data = s3client.get_object(Key=s3key, Bucket=RAW_FILE_BUCKET)
     with tempfile.NamedTemporaryFile() as fd:
         raw_file_contents = data['Body'].read()
         fd.write(raw_file_contents)
         fd.flush()
 
-        if not should_unpack(job_key):
-            log.info("Not unpacking %s", job_key)
-            return [(job_key, raw_file_contents)]
+        if not should_unpack(s3key):
+            log.info("Not unpacking %s", s3key)
+            return [(s3key, raw_file_contents)]
 
         to_extract = archive_contents(fd)
         if not to_extract:
@@ -87,7 +91,7 @@ def unpack(job_key, s3client, redis):
             log.info('Processing %s %s', index, fname)
             file_contents = check_output(['unar', '-o', '-', '-i', fd.name, str(index)])
             log.info('Got %d bytes', len(file_contents))
-            ret.append(fname, file_contents)
+            ret.append((fname, file_contents))
         return ret