forked from werwolfby/monitorrent
-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathengine.py
More file actions
149 lines (122 loc) · 4.32 KB
/
Copy pathengine.py
File metadata and controls
149 lines (122 loc) · 4.32 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
import threading
from datetime import datetime
from sqlalchemy import Column, Integer, DateTime
from db import Base, DBSession
class Logger(object):
def started(self):
pass
def finished(self, finish_time, exception):
pass
def info(self, message):
pass
def failed(self, message):
pass
def downloaded(self, message, torrent):
pass
class Engine(object):
def __init__(self, logger, clients_manager):
"""
:type logger: Logger
:type clients_manager: plugin_managers.ClientsManager
"""
self.log = logger
self.clients_manager = clients_manager
def find_torrent(self, torrent_hash):
return self.clients_manager.find_torrent(torrent_hash)
def add_torrent(self, torrent):
return self.clients_manager.add_torrent(torrent)
def remove_torrent(self, torrent_hash):
return self.clients_manager.remove_torrent(torrent_hash)
class Execute(Base):
__tablename__ = "settings_execute"
id = Column(Integer, primary_key=True)
interval = Column(Integer, nullable=False)
last_execute = Column(DateTime, nullable=True)
class EngineRunner(object):
def __init__(self, logger, trackers_manager, clients_manager):
"""
:type logger: Logger
:type trackers_manager: plugin_managers.TrackersManager
:type clients_manager: plugin_managers.ClientsManager
"""
self.logger = logger
self.trackers_manager = trackers_manager
self.clients_manager = clients_manager
self._execute_lock = threading.RLock()
self._is_executing = False
self.timer = threading.Timer(self.interval, self._run)
self.timer.start()
self._timer_lock = threading.RLock()
@property
def is_executing(self):
with self._execute_lock:
return self._is_executing
@property
def interval(self):
settings_execute = self._get_settings_execute()
return settings_execute.interval
@interval.setter
def interval(self, value):
settings_execute = self._get_settings_execute()
with DBSession() as db:
db.add(settings_execute)
settings_execute.interval = value
db.commit()
self.timer.cancel()
with self._timer_lock:
self.timer = threading.Timer(value, self._run)
self.timer.start()
@property
def last_execute(self):
settings_execute = self._get_settings_execute()
return settings_execute.last_execute
def start(self):
self.timer.start()
def stop(self):
self.timer.cancel()
def execute(self):
caught_exception = None
with self._execute_lock:
if self._is_executing:
return False
self._is_executing = True
try:
self.logger.started()
self.trackers_manager.execute(Engine(self.logger, self.clients_manager))
except Exception as e:
caught_exception = e
finally:
finish_time = datetime.now()
self.logger.finished(finish_time, caught_exception)
self._set_last_execute(finish_time)
with self._execute_lock:
self._is_executing = False
return True
def _run(self):
with self._timer_lock:
old_timer = self.timer
try:
self.execute()
finally:
with self._timer_lock:
# if timer was changed by update interval property
# do not restart the timer
if self.timer == old_timer:
self.timer = threading.Timer(self.interval, self._run)
self.timer.start()
def _set_last_execute(self, value):
settings_execute = self._get_settings_execute()
with DBSession() as db:
db.add(settings_execute)
settings_execute.last_execute = value
db.commit()
@staticmethod
def _get_settings_execute():
with DBSession() as db:
if db.query(Execute).count() == 0:
settings_execute = Execute(interval=7200, last_execute=None)
db.add(settings_execute)
db.commit()
settings_execute = db.query(Execute).first()
db.expunge(settings_execute)
return settings_execute