diff --git a/core/cleanup.py b/core/cleanup.py index a0d0c62..3a82fa8 100644 --- a/core/cleanup.py +++ b/core/cleanup.py @@ -34,6 +34,10 @@ class CleanupTimer: """Stop any threads and prepare for application shutdown""" self._stop_event.set() + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning("Cleanup worker thread did not exit on time and will be killed.") def _run(self): while not self._stop_event.wait(timeout=self._cleanup_interval): diff --git a/core/data_providers.py b/core/data_providers.py index e0a7e5f..06423ee 100644 --- a/core/data_providers.py +++ b/core/data_providers.py @@ -1,5 +1,6 @@ import logging import threading +import time from core.config import config, create_provider_from_config @@ -16,6 +17,7 @@ class DataProviders: self.static_data_providers = [] self.sig_ref_data_providers = [] self.callsign_data_providers = [] + self._startup_timers = [] def setup(self): for entry in config["spot_providers"]: @@ -43,41 +45,58 @@ class DataProviders: def start(self): # Start data providers before spot/alert providers so the lookup data is there already for incoming spots. # Each category is fired off after a small delay to give the rest of Spothole chance to start up. - threading.Timer(5.0, lambda: self.start_providers(self.static_data_providers, "static data")).start() - threading.Timer( - 10.0, - lambda: self.start_providers(self.callsign_data_providers, "callsign data"), - ).start() - threading.Timer(15.0, lambda: self.start_providers(self.spot_providers, "spot")).start() - threading.Timer(20.0, lambda: self.start_providers(self.alert_providers, "alert")).start() - threading.Timer( - 25.0, - lambda: self.start_providers(self.solar_condition_providers, "solar condition"), - ).start() - threading.Timer( - 30.0, - lambda: self.start_providers(self.sig_ref_data_providers, "SIG ref data"), - ).start() + self._startup_timers = [ + threading.Timer(5.0, lambda: self.start_providers(self.static_data_providers, "static data")), + threading.Timer(10.0, lambda: self.start_providers(self.callsign_data_providers, "callsign data")), + threading.Timer(15.0, lambda: self.start_providers(self.spot_providers, "spot")), + threading.Timer(20.0, lambda: self.start_providers(self.alert_providers, "alert")), + threading.Timer( + 25.0, + lambda: self.start_providers(self.solar_condition_providers, "solar condition"), + ), + threading.Timer(30.0, lambda: self.start_providers(self.sig_ref_data_providers, "SIG ref data")), + ] + for t in self._startup_timers: + t.daemon = True + t.start() def stop(self): - for sp in self.spot_providers: - if sp.enabled: - sp.stop() - for ap in self.alert_providers: - if ap.enabled: - ap.stop() - for scp in self.solar_condition_providers: - if scp.enabled: - scp.stop() - for srdp in self.sig_ref_data_providers: - if srdp.enabled: - srdp.stop() - for sdp in self.static_data_providers: - if sdp.enabled: - sdp.stop() - for cdp in self.callsign_data_providers: - if cdp.enabled: - cdp.stop() + # Cancel any startup timers that haven't fired yet + for t in self._startup_timers: + t.cancel() + + # Stop all providers + all_providers = [ + p + for p in ( + self.spot_providers + + self.alert_providers + + self.solar_condition_providers + + self.sig_ref_data_providers + + self.static_data_providers + + self.callsign_data_providers + ) + if p.enabled + ] + if not all_providers: + return + + def stop_provider(p): + try: + p.stop() + except Exception: + logger.exception("Exception stopping provider") + + threads = [threading.Thread(target=stop_provider, args=(p,), daemon=True) for p in all_providers] + for t in threads: + t.start() + + deadline = time.monotonic() + 40 + for t in threads: + t.join(timeout=max(0.0, deadline - time.monotonic())) + still_running = [t for t in threads if t.is_alive()] + if still_running: + logger.warning("Some threads did not stop in time!") # Global object diff --git a/core/live_data_cache.py b/core/live_data_cache.py index 7d13cb7..6fb61e4 100644 --- a/core/live_data_cache.py +++ b/core/live_data_cache.py @@ -23,6 +23,7 @@ class LiveDataCache: self._snapshot_dir = snapshot_dir self._disk_cache = diskcache.Cache(str(snapshot_dir)) self._stop_event = threading.Event() + self._snapshot_thread = None self._load_snapshot() self._start_periodic_snapshot(snapshot_interval_sec) @@ -93,10 +94,16 @@ class LiveDataCache: while not self._stop_event.wait(timeout=interval): self.save_snapshot() - t = threading.Thread(target=loop, name=f"LiveDataCache-Snapshot-{self._snapshot_dir}", daemon=True) - t.start() + self._snapshot_thread = threading.Thread( + target=loop, name=f"LiveDataCache-Snapshot-{self._snapshot_dir}", daemon=True + ) + self._snapshot_thread.start() def close(self): self._stop_event.set() + if self._snapshot_thread: + self._snapshot_thread.join(timeout=15) + if self._snapshot_thread.is_alive(): + logger.warning(f"LiveDataCache snapshot thread for {self._snapshot_dir} did not exit on time.") self.save_snapshot() self._disk_cache.close() diff --git a/core/status_reporter.py b/core/status_reporter.py index 609f938..58c8b36 100644 --- a/core/status_reporter.py +++ b/core/status_reporter.py @@ -1,3 +1,4 @@ +import logging import os from datetime import datetime from threading import Event, Thread @@ -14,6 +15,8 @@ from core.prometheus_metrics_handler import alerts_gauge, memory_use_gauge, spot from telnetserver.telnetserver import TELNET_SERVER from webserver.webserver import WEB_SERVER +logger = logging.getLogger(__name__) + class StatusReporter: """Provides a timed update of the application's status data.""" @@ -33,13 +36,17 @@ class StatusReporter: def start(self): """Start the reporter thread""" - self._thread = Thread(target=self._run, name="StatusReporter") + self._thread = Thread(target=self._run, name="StatusReporter", daemon=True) self._thread.start() def stop(self): """Stop any threads and prepare for application shutdown""" self._stop_event.set() + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning("Status reporter worker thread did not exit on time and will be killed.") def _run(self): """Thread entry point: report immediately on startup, then on each interval until stopped""" diff --git a/providers/alert/http_alert_provider.py b/providers/alert/http_alert_provider.py index 93aab50..83d1c6b 100644 --- a/providers/alert/http_alert_provider.py +++ b/providers/alert/http_alert_provider.py @@ -27,11 +27,15 @@ class HTTPAlertProvider(AlertProvider): # Fire off the polling thread. It will poll immediately on startup, then sleep for poll_interval between # subsequent polls, so start() returns immediately and the application can continue starting. logger.info(f"Set up query of {self.name} alert API every {self._poll_interval!s} seconds.") - self._thread = Thread(target=self._run, name=f"HTTPAlertProvider-{self.name}") + self._thread = Thread(target=self._run, name=f"HTTPAlertProvider-{self.name}", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.name} alert worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/providers/callsigndata/file_download_callsign_data_provider.py b/providers/callsigndata/file_download_callsign_data_provider.py index 7a7a31d..0663ea9 100644 --- a/providers/callsigndata/file_download_callsign_data_provider.py +++ b/providers/callsigndata/file_download_callsign_data_provider.py @@ -32,11 +32,15 @@ class FileDownloadCallsignDataProvider(CallsignDataProvider): # Fire off the polling thread. It will poll immediately on startup, then sleep for poll_interval between # subsequent polls, so start() returns immediately and the application can continue starting. logger.info(f"Set up query of {self.name} callsign reference data every {self._poll_interval!s} days.") - self._thread = Thread(target=self._run, name=f"FileDownloadCallsignDataProvider-{self.name}") + self._thread = Thread(target=self._run, name=f"FileDownloadCallsignDataProvider-{self.name}", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.name} callsign data worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/providers/sigrefdata/file_download_sig_ref_data_provider.py b/providers/sigrefdata/file_download_sig_ref_data_provider.py index 92369d7..d1e243b 100644 --- a/providers/sigrefdata/file_download_sig_ref_data_provider.py +++ b/providers/sigrefdata/file_download_sig_ref_data_provider.py @@ -1,6 +1,6 @@ import logging from datetime import datetime -from threading import Event, Thread +from threading import Thread import pytz from requests.exceptions import ConnectionError, ConnectTimeout, ReadTimeout @@ -21,19 +21,21 @@ class FileDownloadSIGRefDataProvider(SIGRefDataProvider): self._url = url self._poll_interval = poll_interval self._thread = None - self._stop_event = Event() self._url_data_cache = URLDataCache(f"sigrefdata_{sig_name}") def start(self): # Fire off the polling thread. It will poll immediately on startup, then sleep for poll_interval between # subsequent polls, so start() returns immediately and the application can continue starting. logger.info(f"Set up query of {self.sig_name} SIG ref data every {self._poll_interval!s} days.") - self._thread = Thread(target=self._run, name=f"FileDownloadSIGRefDataProvider-{self.sig_name}") + self._thread = Thread(target=self._run, name=f"FileDownloadSIGRefDataProvider-{self.sig_name}", daemon=True) self._thread.start() def stop(self): super().stop() - self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.sig_name} SIG ref data worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/providers/sigrefdata/sig_ref_data_provider.py b/providers/sigrefdata/sig_ref_data_provider.py index 7c98345..bf128a9 100644 --- a/providers/sigrefdata/sig_ref_data_provider.py +++ b/providers/sigrefdata/sig_ref_data_provider.py @@ -1,5 +1,6 @@ import logging from datetime import datetime +from threading import Event import pytz @@ -19,7 +20,7 @@ class SIGRefDataProvider: self.last_update_time = datetime.min.replace(tzinfo=pytz.UTC) self.status = "Not Started" if self.enabled else "Disabled" self.reference_count = 0 - self._stop = False + self._stop_event = Event() def start(self): """Start the provider. This should return immediately after spawning threads to access the remote resources""" @@ -30,7 +31,7 @@ class SIGRefDataProvider: """Stop any threads and prepare for application shutdown. Subclasses should implement this method and call super().""" - self._stop = True + self._stop_event.set() def _add_data(self, new_data): """Add all the provided reference data objects to the data store.""" @@ -43,7 +44,7 @@ class SIGRefDataProvider: # For the big data sources, loading will take a few minutes. If we want to shut down the software neatly # within the first few minutes of startup, we need a way to abort this expensive process of filling up the # disk cache. - if self._stop: + if self._stop_event.is_set(): break self.reference_count = len(new_data) diff --git a/providers/solarconditions/giroionosonde.py b/providers/solarconditions/giroionosonde.py index 023ca23..4067920 100644 --- a/providers/solarconditions/giroionosonde.py +++ b/providers/solarconditions/giroionosonde.py @@ -67,11 +67,15 @@ class GIROIonosonde(SolarConditionsProvider): def start(self): logger.info(f"Set up query of GIRO ionosonde data API every {POLL_INTERVAL} seconds.") - self._thread = Thread(target=self._run, name="GIROIonosondeDataProvider") + self._thread = Thread(target=self._run, name="GIROIonosondeDataProvider", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning("GIRO ionosonde worker thread did not exit on time and will be killed.") def _run(self): # Real interval at which we poll is the "once per hour" divided by the number of stations, so each one gets diff --git a/providers/solarconditions/http_solar_conditions_provider.py b/providers/solarconditions/http_solar_conditions_provider.py index 86229dc..a5feeb9 100644 --- a/providers/solarconditions/http_solar_conditions_provider.py +++ b/providers/solarconditions/http_solar_conditions_provider.py @@ -25,11 +25,15 @@ class HTTPSolarConditionsProvider(SolarConditionsProvider): def start(self): logger.info(f"Set up query of {self.name} solar conditions API every {self._poll_interval!s} seconds.") - self._thread = Thread(target=self._run, name=f"HTTPSolarConditionsProvider-{self.name}") + self._thread = Thread(target=self._run, name=f"HTTPSolarConditionsProvider-{self.name}", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.name} solar conditions worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/providers/solarconditions/kc2gprop.py b/providers/solarconditions/kc2gprop.py index f80f423..8f9cd81 100644 --- a/providers/solarconditions/kc2gprop.py +++ b/providers/solarconditions/kc2gprop.py @@ -32,11 +32,15 @@ class KC2GProp(SolarConditionsProvider): def start(self): logger.info(f"Set up query of KC2G ionosonde data API every {POLL_INTERVAL} seconds.") - self._thread = Thread(target=self._run, name="KC2GPropProvider") + self._thread = Thread(target=self._run, name="KC2GPropProvider", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning("KC2G ionosonde worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/providers/spot/aprsis.py b/providers/spot/aprsis.py index 9442881..3977a60 100644 --- a/providers/spot/aprsis.py +++ b/providers/spot/aprsis.py @@ -17,17 +17,16 @@ class APRSIS(SpotProvider): def __init__(self, provider_config): super().__init__("APRS-IS", provider_config) - self._thread = Thread(target=self._run, name="APRSISSpotProvider") - self._thread.daemon = True + self._thread = None self._aprsis = None - self._running = True self._stop_event = Event() def start(self): + self._thread = Thread(target=self._run, name="APRSISSpotProvider", daemon=True) self._thread.start() def _run(self): - while self._running: + while not self._stop_event.is_set(): try: self._aprsis = aprslib.IS(SERVER_OWNER_CALLSIGN) self.status = "Connecting" @@ -37,20 +36,22 @@ class APRSIS(SpotProvider): self._aprsis.consumer(self._handle, immortal=True) except Exception: - if self._running: + if not self._stop_event.is_set(): self.status = "Error" logger.exception("Exception in APRS-IS provider") - if self._running: + if not self._stop_event.is_set(): self._stop_event.wait(timeout=5) def stop(self): - self._running = False self.status = "Shutting down" self._stop_event.set() if self._aprsis: self._aprsis.close() - self._thread.join() + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning("APRS-IS worker thread did not exit on time and will be killed.") def _handle(self, data): try: diff --git a/providers/spot/dxcluster.py b/providers/spot/dxcluster.py index 5cbe78f..eb537df 100644 --- a/providers/spot/dxcluster.py +++ b/providers/spot/dxcluster.py @@ -1,8 +1,7 @@ import logging import re from datetime import datetime -from threading import Thread -from time import sleep +from threading import Event, Thread import pytz import telnetlib3 @@ -41,23 +40,26 @@ class DXCluster(SpotProvider): self._LINE_PATTERN_ALLOW_RBN if self._allow_rbn_spots else self._LINE_PATTERN_EXCLUDE_RBN ) self._telnet = None - self._thread = Thread(target=self._handle, name=f"DXClusterSpotProvider-{self.name}") - self._thread.daemon = True - self._running = True + self._thread = None + self._stop_event = Event() def start(self): + self._thread = Thread(target=self._handle, name=f"DXClusterSpotProvider-{self.name}", daemon=True) self._thread.start() def stop(self): - self._running = False + self._stop_event.set() if self._telnet: self._telnet.close() - self._thread.join() + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning(f"DX Cluster {self._hostname} worker thread did not exit on time and will be killed.") def _handle(self): - while self._running: + while not self._stop_event.is_set(): connected = False - while not connected and self._running: + while not connected and not self._stop_event.is_set(): try: self.status = "Connecting" logger.info(f"DX Cluster {self._hostname} connecting...") @@ -69,14 +71,14 @@ class DXCluster(SpotProvider): except ConnectionRefusedError: self.status = "Error" logger.warning(f"Connection refused to DX cluster {self._hostname}") - sleep(300) + self._stop_event.wait(timeout=300) except Exception: self.status = "Error" logger.exception(f"Exception while connecting to DX Cluster Provider ({self._hostname}).") - sleep(5) + self._stop_event.wait(timeout=5) self.status = "Waiting for Data" - while connected and self._running: + while connected and not self._stop_event.is_set(): try: # Check new telnet info against regular expression telnet_output = self._telnet.read_until("\n".encode("latin-1")) @@ -106,19 +108,19 @@ class DXCluster(SpotProvider): except EOFError: connected = False - if self._running: + if not self._stop_event.is_set(): self.status = "Restarting" logger.warning(f"Disconnected from DX Cluster {self._hostname}. Reconnecting...") - sleep(5) + self._stop_event.wait(timeout=5) else: logger.info(f"DX Cluster {self._hostname} shutting down...") self.status = "Shutting down" except Exception: connected = False - if self._running: + if not self._stop_event.is_set(): self.status = "Error" logger.exception(f"Exception in DX Cluster Provider ({self._hostname})") - sleep(5) + self._stop_event.wait(timeout=5) else: logger.info(f"DX Cluster {self._hostname} shutting down...") self.status = "Shutting down" diff --git a/providers/spot/http_spot_provider.py b/providers/spot/http_spot_provider.py index 14b4fa4..124742a 100644 --- a/providers/spot/http_spot_provider.py +++ b/providers/spot/http_spot_provider.py @@ -28,12 +28,16 @@ class HTTPSpotProvider(SpotProvider): # Fire off the polling thread. It will poll immediately on startup, then sleep for poll_interval between # subsequent polls, so start() returns immediately and the application can continue starting. logger.info(f"Set up query of {self.name} spot API every {self._poll_interval!s} seconds.") - self._thread = Thread(target=self._run, name=f"HTTPSpotProvider-{self.name}") + self._thread = Thread(target=self._run, name=f"HTTPSpotProvider-{self.name}", daemon=True) self._thread.start() def stop(self): self._stop_event.set() self._wakeup_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.name} spot worker thread did not exit on time and will be killed.") def force_poll(self): """Trigger an immediate poll without waiting for the normal interval.""" diff --git a/providers/spot/rbn.py b/providers/spot/rbn.py index 6b26238..7923b9e 100644 --- a/providers/spot/rbn.py +++ b/providers/spot/rbn.py @@ -1,8 +1,7 @@ import logging import re from datetime import datetime -from threading import Thread -from time import sleep +from threading import Event, Thread import pytz import telnetlib3 @@ -30,23 +29,26 @@ class RBN(SpotProvider): super().__init__(name, provider_config) self._port = provider_config["port"] self._telnet = None - self._thread = Thread(target=self._handle, name=f"RBNSpotProvider-{self.name}") - self._thread.daemon = True - self._running = True + self._thread = None + self._stop_event = Event() def start(self): + self._thread = Thread(target=self._handle, name=f"RBNSpotProvider-{self.name}", daemon=True) self._thread.start() def stop(self): - self._running = False + self._stop_event.set() if self._telnet: self._telnet.close() - self._thread.join() + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning(f"RBN (port {self._port!s}) worker thread did not exit on time and will be killed.") def _handle(self): - while self._running: + while not self._stop_event.is_set(): connected = False - while not connected and self._running: + while not connected and not self._stop_event.is_set(): try: self.status = "Connecting" logger.info(f"RBN port {self._port!s} connecting...") @@ -58,10 +60,10 @@ class RBN(SpotProvider): except Exception: self.status = "Error" logger.exception(f"Exception while connecting to RBN (port {self._port!s}).") - sleep(5) + self._stop_event.wait(timeout=5) self.status = "Waiting for Data" - while connected and self._running: + while connected and not self._stop_event.is_set(): try: # Check new telnet info against regular expression telnet_output = self._telnet.read_until("\n".encode("latin-1")) @@ -91,19 +93,19 @@ class RBN(SpotProvider): except EOFError: connected = False - if self._running: + if not self._stop_event.is_set(): self.status = "Restarting" logger.warning(f"Disconnected from RBN provider (port {self._port!s}). Reconnecting...") - sleep(5) + self._stop_event.wait(timeout=5) else: logger.info(f"RBN provider (port {self._port!s}) shutting down...") self.status = "Shutting down" except Exception: connected = False - if self._running: + if not self._stop_event.is_set(): self.status = "Error" logger.exception(f"Exception in RBN provider (port {self._port!s})") - sleep(5) + self._stop_event.wait(timeout=5) else: logger.info(f"RBN provider (port {self._port!s}) shutting down...") self.status = "Shutting down" diff --git a/providers/spot/websocket_spot_provider.py b/providers/spot/websocket_spot_provider.py index f5b2ebf..5c5af61 100644 --- a/providers/spot/websocket_spot_provider.py +++ b/providers/spot/websocket_spot_provider.py @@ -1,7 +1,6 @@ import logging from datetime import datetime -from threading import Thread -from time import sleep +from threading import Event, Thread import pytz from websocket import create_connection @@ -20,22 +19,24 @@ class WebsocketSpotProvider(SpotProvider): self._url = url self._ws = None self._thread = None - self._stopped = False + self._stop_event = Event() self._last_event_id = None def start(self): logger.info(f"Set up websocket connection to {self.name} spot API.") - self._stopped = False + self._stop_event.clear() self._thread = Thread(target=self._run, name=f"WebsocketSpotProvider-{self.name}") self._thread.daemon = True self._thread.start() def stop(self): - self._stopped = True + self._stop_event.set() if self._ws: self._ws.close() if self._thread: - self._thread.join() + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning(f"{self.name} websocket worker thread did not exit on time and will be killed.") def _on_open(self): self.status = "Waiting for Data" @@ -44,7 +45,7 @@ class WebsocketSpotProvider(SpotProvider): self.status = "Connecting" def _run(self): - while not self._stopped: + while not self._stop_event.is_set(): try: logger.debug(f"Connecting to {self.name} spot API...") self.status = "Connecting" @@ -53,7 +54,7 @@ class WebsocketSpotProvider(SpotProvider): # Keep reading from this same connection until it drops or we're asked to stop, rather than # reconnecting for every message. - while not self._stopped: + while not self._stop_event.is_set(): data = self._ws.recv() if not data: break @@ -82,8 +83,8 @@ class WebsocketSpotProvider(SpotProvider): # No problem, we were getting rid of this object anyway. pass self._ws = None - if not self._stopped: - sleep(5) # Wait before trying to reconnect + if not self._stop_event.is_set(): + self._stop_event.wait(timeout=5) # Wait before trying to reconnect def _ws_message_to_spot(self, b): """Convert a WS message received from the API into a spot. The exact message data (in bytes) is provided here so the diff --git a/providers/staticdata/file_download_static_data_provider.py b/providers/staticdata/file_download_static_data_provider.py index 44a26a3..6b733f6 100644 --- a/providers/staticdata/file_download_static_data_provider.py +++ b/providers/staticdata/file_download_static_data_provider.py @@ -29,11 +29,15 @@ class FileDownloadStaticDataProvider(StaticDataProvider): # Fire off the polling thread. It will poll immediately on startup, then sleep for poll_interval between # subsequent polls, so start() returns immediately and the application can continue starting. logger.info(f"Set up query of {self.name} static reference data every {self._poll_interval!s} days.") - self._thread = Thread(target=self._run, name=f"FileDownloadStaticDataProvider-{self.name}") + self._thread = Thread(target=self._run, name=f"FileDownloadStaticDataProvider-{self.name}", daemon=True) self._thread.start() def stop(self): self._stop_event.set() + if self._thread: + self._thread.join(timeout=35) + if self._thread.is_alive(): + logger.warning(f"{self.name} static data worker thread did not exit on time and will be killed.") def _run(self): while True: diff --git a/spothole.py b/spothole.py index b9835b5..bcc8066 100644 --- a/spothole.py +++ b/spothole.py @@ -15,17 +15,25 @@ from webserver.webserver import WEB_SERVER logger = logging.getLogger(__name__) +_shutdown_in_progress = False + def shutdown(_signum=None, _frame=None): """Shutdown function""" + # Check if this is the second time a shutdown was asked for, if so immediately kill the program. + global _shutdown_in_progress + if _shutdown_in_progress: + os._exit(1) + _shutdown_in_progress = True + logger.info("Stopping program...") WEB_SERVER.stop() TELNET_SERVER.stop() DATA_PROVIDERS.stop() CLEANUP_TIMER.stop() DATA_STORE.close() - os._exit(0) + logger.info("Stopped.") # Main function @@ -43,8 +51,9 @@ if __name__ == "__main__": logger.info("Starting...") logger.info(f"This is Spothole version {SOFTWARE_VERSION}. This instance is run by {SERVER_OWNER_CALLSIGN}.") - # Shut down gracefully on SIGINT + # Shut down gracefully on SIGINT or SIGTERM signal.signal(signal.SIGINT, shutdown) + signal.signal(signal.SIGTERM, shutdown) # Set up data store DATA_STORE.setup() @@ -59,14 +68,16 @@ if __name__ == "__main__": status_reporter = StatusReporter(run_interval=5) status_reporter.start() - # Set up the web server - WEB_SERVER.setup() - # Run the telnet server if TELNET_SERVER_ENABLED: TELNET_SERVER.start(port=TELNET_SERVER_PORT) - # Run the web server. This is the blocking call that keeps the application running in the main thread, so this must - # be the last thing we do. web_server.stop() triggers an await condition in the web server which finishes the main - # thread. + # Set up the web server + WEB_SERVER.setup() + + # Run the web server WEB_SERVER.start() + + # Block the main thread until a termination signal arrives and shutdown() is running. + while not _shutdown_in_progress: + signal.pause() diff --git a/telnetserver/telnetserver.py b/telnetserver/telnetserver.py index dfb3393..57996ea 100644 --- a/telnetserver/telnetserver.py +++ b/telnetserver/telnetserver.py @@ -49,6 +49,7 @@ class TelnetServer: self._running = False self._clients = set() self._loop = None + self._thread = None self._shutdown_event = asyncio.Event() def start(self, port=7373): @@ -58,8 +59,10 @@ class TelnetServer: # Start the telnet server. asyncio.run() needs a coroutine, and threading.Thread needs a plain callable, so # hand Thread the bridge between the two directly rather than writing a one-line wrapper method for it. - t = threading.Thread(target=asyncio.run, args=(self._start_internal(),), name="TelnetServer", daemon=True) - t.start() + self._thread = threading.Thread( + target=asyncio.run, args=(self._start_internal(),), name="TelnetServer", daemon=True + ) + self._thread.start() logger.debug("Telnet server background thread spawned") # Listen for new spots and alerts being added to the cache, so we can notify SSE clients immediately @@ -135,6 +138,10 @@ class TelnetServer: if self._loop and self._loop.is_running(): logger.debug("Stopping telnet server...") self._loop.call_soon_threadsafe(self._shutdown_event.set) + if self._thread: + self._thread.join(timeout=15) + if self._thread.is_alive(): + logger.warning("Telnet server background thread did not exit on time and will be killed.") @property def client_count(self) -> int: diff --git a/templates/add_spot.html b/templates/add_spot.html index 8be35a2..8c82ead 100644 --- a/templates/add_spot.html +++ b/templates/add_spot.html @@ -77,7 +77,7 @@ - + diff --git a/templates/alerts.html b/templates/alerts.html index 27b4c38..8f273b1 100644 --- a/templates/alerts.html +++ b/templates/alerts.html @@ -83,7 +83,7 @@ - + diff --git a/templates/bands.html b/templates/bands.html index 3ed1684..93a9b9a 100644 --- a/templates/bands.html +++ b/templates/bands.html @@ -76,8 +76,8 @@ - - + + diff --git a/templates/base.html b/templates/base.html index 7a35986..5ed9f5b 100644 --- a/templates/base.html +++ b/templates/base.html @@ -1,6 +1,6 @@ {% extends "skeleton.html" %} {% block head_extra %} - + @@ -16,10 +16,10 @@ window.fetchEventSource = fetchEventSource; - - - - + + + + {% end %} {% block body %}
diff --git a/templates/conditions.html b/templates/conditions.html index 58761e9..ac538de 100644 --- a/templates/conditions.html +++ b/templates/conditions.html @@ -284,7 +284,7 @@
- + diff --git a/templates/map.html b/templates/map.html index 810030a..cb3f0a4 100644 --- a/templates/map.html +++ b/templates/map.html @@ -113,8 +113,8 @@ const CARTODB_API_KEY = "{{ web_ui_options.get('cartodb_api_key', '') }}"; - - + + diff --git a/templates/spots.html b/templates/spots.html index 1984dfb..d7da28c 100644 --- a/templates/spots.html +++ b/templates/spots.html @@ -125,8 +125,8 @@ - - + + diff --git a/templates/status.html b/templates/status.html index c813ee2..335fb10 100644 --- a/templates/status.html +++ b/templates/status.html @@ -96,7 +96,7 @@ - +