|
|
@@ -1,36 +1,44 @@
|
|
|
+import os
|
|
|
+import random
|
|
|
+import string
|
|
|
import threading
|
|
|
-import queue
|
|
|
import logging
|
|
|
-import functools
|
|
|
+import time
|
|
|
|
|
|
-q = queue.Queue()
|
|
|
+from config import UPLOAD_DIR
|
|
|
|
|
|
|
|
|
-def worker():
|
|
|
- while True:
|
|
|
+def wait_and_execute(job, callback):
|
|
|
+ def worker():
|
|
|
+ while True:
|
|
|
+ count = 0
|
|
|
+ for i in os.listdir(UPLOAD_DIR):
|
|
|
+ if i.endswith(".running"):
|
|
|
+ count += 1
|
|
|
+ if count < 2:
|
|
|
+ break
|
|
|
+ time.sleep(random.randint(1, 2))
|
|
|
+
|
|
|
+ name = ''.join(random.sample(string.ascii_letters + string.digits, 8))
|
|
|
+ with open(UPLOAD_DIR + "/" + name + ".running", "w") as f:
|
|
|
+ pass
|
|
|
+
|
|
|
+ logging.info("start job: " + name)
|
|
|
+
|
|
|
try:
|
|
|
- task = q.get()
|
|
|
- logging.debug(f"Working on {task}")
|
|
|
- job, callback = task
|
|
|
job(callback)
|
|
|
- logging.debug(f"Working on {task} done.")
|
|
|
- q.task_done()
|
|
|
except Exception as e:
|
|
|
- logging.debug("error:", exc_info=e)
|
|
|
+ logging.debug("Worker error: ", exc_info=e)
|
|
|
raise e
|
|
|
|
|
|
+ try:
|
|
|
+ os.remove(UPLOAD_DIR + "/" + name + ".running")
|
|
|
+ logging.info("finish job: " + name)
|
|
|
+ except:
|
|
|
+ pass
|
|
|
|
|
|
-def init_workers(count=2):
|
|
|
- logging.info("Initializing workers...")
|
|
|
- global q
|
|
|
- q = queue.Queue()
|
|
|
- for _ in range(count):
|
|
|
- threading.Thread(target=worker).start()
|
|
|
-
|
|
|
-
|
|
|
-class Pool:
|
|
|
- def submit(self, job, callback):
|
|
|
- q.put((job, callback))
|
|
|
+ return worker
|
|
|
|
|
|
|
|
|
-pool = Pool()
|
|
|
+def submit(job, callback):
|
|
|
+ threading.Thread(target=wait_and_execute(job, callback)).start()
|