version 0.3.0, rev187, Trusted authorization sites support, --publish option on signing, cryptSign command line option, OpenSSL enabled on OSX, Crypto verify allows list of valid addresses, Option for version 2 json DB tables, DbCursor SELECT parameters bugfix, Add peer to site on ListModified, Download blind includes when new site added, Publish command better messages, Multi-threaded announce, New http Torrent trackers, Wait for dbschema.json on query, Handle json import errors, More compact writeJson storage command, Testcase for signing and verifying, Workaround to make non target=_top links work, More clean UiWebsocket command route, Send cert_user_id on siteinfo, Notify other local clients on local file modify, Option to wait for file download before sql query, File rules websocket API command, Cert add and select, set websocket API command, Put focus on innerframe, innerloaded wrapper api command to add hashtag, Allow more file error on big sites, Keep worker running after stuked on done task, New more stable openSSL layer that works on OSX, Noparallel parameter bugfix, RateLimit allowed again interval bugfix, Updater skips non-writeable files, Try to close openssl dll before update
This commit is contained in:
parent
c874726aba
commit
7e4f6bd38e
33 changed files with 1716 additions and 595 deletions
|
@ -62,7 +62,7 @@ class Worker:
|
|||
task["failed"].append(self.peer)
|
||||
self.task = None
|
||||
self.peer.hash_failed += 1
|
||||
if self.peer.hash_failed >= 3: # Broken peer
|
||||
if self.peer.hash_failed >= max(len(self.manager.tasks), 3): # More fails than tasks number but atleast 3: Broken peer
|
||||
break
|
||||
task["workers_num"] -= 1
|
||||
time.sleep(1)
|
||||
|
@ -77,9 +77,17 @@ class Worker:
|
|||
self.thread = gevent.spawn(self.downloader)
|
||||
|
||||
|
||||
# Skip current task
|
||||
def skip(self):
|
||||
self.manager.log.debug("%s: Force skipping" % self.key)
|
||||
if self.thread:
|
||||
self.thread.kill(exception=Debug.Notify("Worker stopped"))
|
||||
self.start()
|
||||
|
||||
|
||||
# Force stop the worker
|
||||
def stop(self):
|
||||
self.manager.log.debug("%s: Force stopping, thread" % self.key)
|
||||
self.manager.log.debug("%s: Force stopping" % self.key)
|
||||
self.running = False
|
||||
if self.thread:
|
||||
self.thread.kill(exception=Debug.Notify("Worker stopped"))
|
||||
|
|
|
@ -32,18 +32,23 @@ class WorkerManager:
|
|||
|
||||
# Clean up workers
|
||||
for worker in self.workers.values():
|
||||
if worker.task and worker.task["done"]: worker.stop() # Stop workers with task done
|
||||
if worker.task and worker.task["done"]: worker.skip() # Stop workers with task done
|
||||
|
||||
if not self.tasks: continue
|
||||
|
||||
tasks = self.tasks[:] # Copy it so removing elements wont cause any problem
|
||||
for task in tasks:
|
||||
if (task["time_started"] and time.time() >= task["time_started"]+60) or (time.time() >= task["time_added"]+60 and not self.workers): # Task taking too long time, or no peer after 60sec kill it
|
||||
self.log.debug("Timeout, Cleaning up task: %s" % task)
|
||||
# Clean up workers
|
||||
if task["time_started"] and time.time() >= task["time_started"]+60: # Task taking too long time, skip it
|
||||
self.log.debug("Timeout, Skipping: %s" % task)
|
||||
# Skip to next file workers
|
||||
workers = self.findWorkers(task)
|
||||
for worker in workers:
|
||||
worker.stop()
|
||||
if workers:
|
||||
for worker in workers:
|
||||
worker.skip()
|
||||
else:
|
||||
self.failTask(task)
|
||||
elif time.time() >= task["time_added"]+60 and not self.workers: # No workers left
|
||||
self.log.debug("Timeout, Cleanup task: %s" % task)
|
||||
# Remove task
|
||||
self.failTask(task)
|
||||
|
||||
|
@ -178,12 +183,13 @@ class WorkerManager:
|
|||
|
||||
# Mark a task failed
|
||||
def failTask(self, task):
|
||||
task["done"] = True
|
||||
self.tasks.remove(task) # Remove from queue
|
||||
self.site.onFileFail(task["inner_path"])
|
||||
task["evt"].set(False)
|
||||
if not self.tasks:
|
||||
self.started_task_num = 0
|
||||
if task in self.tasks:
|
||||
task["done"] = True
|
||||
self.tasks.remove(task) # Remove from queue
|
||||
self.site.onFileFail(task["inner_path"])
|
||||
task["evt"].set(False)
|
||||
if not self.tasks:
|
||||
self.started_task_num = 0
|
||||
|
||||
|
||||
# Mark a task done
|
||||
|
|
Loading…
Add table
Add a link
Reference in a new issue