-
Notifications
You must be signed in to change notification settings - Fork 32
/
cleaner.py
executable file
·252 lines (206 loc) · 8.77 KB
/
cleaner.py
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
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
"""
cleaner.py
==========
The cleaner service of mercure. Responsible for deleting processed data after
retention time has passed and if it is offpeak time. Offpeak is the time
period when the cleaning has to be done, because cleaning I/O should be kept
to minimum when receiving and sending exams.
"""
# Standard python includes
import asyncio
import logging
import os
import signal
import sys
import time
from datetime import timedelta, datetime
from datetime import time as _time
from pathlib import Path
from shutil import rmtree, disk_usage
import graphyte
import hupper
# App-specific includes
import common.config as config
import common.helper as helper
import common.monitor as monitor
from common.monitor import task_event
from common.constants import mercure_defs
import common.influxdb
import common.notification as notification
# Setup daiquiri logger
logger = config.get_logger()
main_loop = None # type: helper.AsyncTimer # type: ignore
async def terminate_process(signalNumber, frame) -> None:
"""Triggers the shutdown of the service."""
helper.g_log("events.shutdown", 1)
logger.info("Shutdown requested")
monitor.send_event(monitor.m_events.SHUTDOWN_REQUEST, monitor.severity.INFO)
# Note: main_loop can be read here because it has been declared as global variable
if "main_loop" in globals() and main_loop.is_running:
main_loop.stop()
helper.trigger_terminate()
def clean() -> None:
"""Main entry function."""
if helper.is_terminated():
return
helper.g_log("events.run", 1)
try:
config.read_config()
except Exception:
logger.warning( # handle_error
"Unable to read configuration. Skipping processing.",
None,
event_type=monitor.m_events.CONFIG_UPDATE,
)
return
## Emergency cleaning procedure: Check if server is running out of disk space. If so, clean images right away
# Get the percentage of disk usage that should trigger the emergency cleaning
emergency_clean_trigger: float = config.mercure.emergency_clean_percentage / 100.0
# Check if the success and discard folder are stored on the same volume
success_folder = config.mercure.success_folder
discard_folder = config.mercure.discard_folder
success_folder_partition = os.stat(success_folder).st_dev
discard_folder_partition = os.stat(discard_folder).st_dev
# For emergency cleaning need to take into account if success and discard
# folders are on the same volume or not.
emergency_retention = timedelta(0)
folders_to_clear = [success_folder, discard_folder]
if success_folder_partition == discard_folder_partition:
(total, used, _) = disk_usage(success_folder)
bytes_to_clear = int(max(used - total * emergency_clean_trigger, 0))
if bytes_to_clear > 0:
for folder in folders_to_clear:
# Need to delete all scan data in the both folders to urgently clean up the space.
clean_dir(folder, emergency_retention)
monitor.send_event(
monitor.m_events.PROCESSING,
monitor.severity.WARNING,
f"Disk is almost full. Emergency cleaning of the {success_folder} and {discard_folder} folders. Consider adjusting retention period.",
)
else:
bytes_to_clear = 0
for folder in folders_to_clear:
(total, used, _) = disk_usage(folder)
bytes_to_clear = int(max(used - total * emergency_clean_trigger, 0))
if bytes_to_clear > 0:
# Need to delete all scan data in the folder to urgently clean up the space.
clean_dir(folder, emergency_retention)
monitor.send_event(
monitor.m_events.PROCESSING,
monitor.severity.WARNING,
f"Disk is almost full. Emergency cleaning of the {folder} folder. Consider adjusting retention period.",
)
## Regular cleaning procedure
if helper._is_offpeak(
config.mercure.offpeak_start,
config.mercure.offpeak_end,
datetime.now().time(),
):
retention = timedelta(seconds=config.mercure.retention)
clean_dir(success_folder, retention)
clean_dir(discard_folder, retention)
def clean_dir(folder, retention) -> None:
"""
Cleans items from the given folder that have exceeded the retention time, starting with the oldest items
"""
candidates = [
(f, f.stat().st_mtime)
for f in Path(folder).iterdir()
if f.is_dir() and retention < timedelta(seconds=(time.time() - f.stat().st_mtime))
]
for entry in candidates:
delete_folder(entry)
def delete_folder(entry) -> None:
"""Deletes given folder."""
delete_path = entry[0]
series_uid = find_series_uid(delete_path)
try:
rmtree(delete_path)
logger.info(f"Deleted folder {delete_path} from {series_uid}")
monitor.send_task_event(task_event.CLEAN, Path(delete_path).stem, 0, delete_path, "Deleted folder")
except Exception as e:
logger.error(
f"Unable to delete folder {delete_path}",
Path(delete_path).stem,
target=delete_path,
) # handle_error
def find_series_uid(work_dir) -> str:
"""
Finds series uid which is always part before the '#'-sign in filename.
"""
to_be_deleted_dir = Path(work_dir)
for entry in to_be_deleted_dir.iterdir():
if "#" in entry.name:
return entry.name.split(mercure_defs.SEPARATOR)[0]
return "series_uid-not-found"
return "series_uid-not-found"
def exit_cleaner(args) -> None:
"""Stop the asyncio event loop."""
helper.loop.call_soon_threadsafe(helper.loop.stop)
def main(args=sys.argv[1:]) -> None:
if "--reload" in args or os.getenv("MERCURE_ENV", "PROD").lower() == "dev":
# start_reloader will only return in a monitored subprocess
reloader = hupper.start_reloader("cleaner.main")
import logging
logging.getLogger("watchdog").setLevel(logging.WARNING)
logger.info("")
logger.info(f"mercure DICOM Cleaner ver {mercure_defs.VERSION}")
logger.info("--------------------------------------------")
logger.info("")
# Register system signals to be caught
signals = (signal.SIGTERM, signal.SIGINT)
for s in signals:
helper.loop.add_signal_handler(s, lambda s=s: asyncio.create_task(terminate_process(s, helper.loop)))
instance_name = "main"
if len(sys.argv) > 1:
instance_name = sys.argv[1]
try:
config.read_config()
except Exception:
logger.exception("Cannot start service. Going down.")
sys.exit(1)
appliance_name = config.mercure.appliance_name
logger.info(f"Appliance name = {appliance_name}")
logger.info(f"Instance name = {instance_name}")
logger.info(f"Instance PID = {os.getpid()}")
logger.info(sys.version)
notification.setup()
monitor.configure("cleaner", instance_name, config.mercure.bookkeeper)
monitor.send_event(monitor.m_events.BOOT, monitor.severity.INFO, f"PID = {os.getpid()}")
if len(config.mercure.graphite_ip) > 0:
logger.info(f"Sending events to graphite server: {config.mercure.graphite_ip}")
graphite_prefix = "mercure." + appliance_name + ".cleaner." + instance_name
graphyte.init(
config.mercure.graphite_ip,
config.mercure.graphite_port,
prefix=graphite_prefix,
)
if len(config.mercure.influxdb_host) > 0:
logger.info(f"Sending events to influxdb server: {config.mercure.influxdb_host}")
common.influxdb.init(
config.mercure.influxdb_host,
config.mercure.influxdb_token,
config.mercure.influxdb_org,
config.mercure.influxdb_bucket,
"mercure." + appliance_name + ".cleaner." + instance_name
)
global main_loop
main_loop = helper.AsyncTimer(config.mercure.cleaner_scan_interval, clean)
main_loop.start()
helper.g_log("events.boot", 1)
try:
# Start the asyncio event loop for asynchronous function calls
main_loop.run_until_complete(helper.loop)
# Process will exit here once the asyncio loop has been stopped
monitor.send_event(monitor.m_events.SHUTDOWN, monitor.severity.INFO)
except Exception as e:
# Process will exit here once the asyncio loop has been stopped
monitor.send_event(monitor.m_events.SHUTDOWN, monitor.severity.ERROR, str(e))
finally:
# Finish all asyncio tasks that might be still pending
remaining_tasks = helper.asyncio.all_tasks(helper.loop) # type: ignore[attr-defined]
if remaining_tasks:
helper.loop.run_until_complete(helper.asyncio.gather(*remaining_tasks))
logger.info("Going down now")
if __name__ == "__main__":
main()