-
Notifications
You must be signed in to change notification settings - Fork 614
Expand file tree
/
Copy pathscheduler_worker.py
More file actions
51 lines (39 loc) · 1.26 KB
/
Copy pathscheduler_worker.py
File metadata and controls
51 lines (39 loc) · 1.26 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
import asyncio
import signal
from config import runtime_settings
from role import Role
runtime_settings.role = Role.SCHEDULER
from app import create_app
from app.lifecycle import lifespan
from app.nats import is_nats_enabled
from app.utils.logger import get_logger
logger = get_logger("scheduler-worker")
app = create_app()
async def main():
if not is_nats_enabled():
logger.warning(
"NATS is disabled; notification dispatching will only work when the scheduler shares a process with the API."
)
stop_event = asyncio.Event()
def handle_signal():
stop_event.set()
loop = asyncio.get_running_loop()
loop.add_signal_handler(signal.SIGINT, handle_signal)
loop.add_signal_handler(signal.SIGTERM, handle_signal)
async with lifespan(app):
try:
logger.info("Scheduler worker started...")
await stop_event.wait()
except asyncio.CancelledError:
pass
finally:
logger.info("Scheduler worker shutting down...")
if __name__ == "__main__":
try:
if hasattr(asyncio, "run"):
asyncio.run(main())
else:
loop = asyncio.get_event_loop()
loop.run_until_complete(main())
except KeyboardInterrupt:
pass