diff --git a/providers/spot/dxcluster.py b/providers/spot/dxcluster.py index ed89d57..911c183 100644 --- a/providers/spot/dxcluster.py +++ b/providers/spot/dxcluster.py @@ -1,21 +1,17 @@ import logging import re -import socket from datetime import datetime -from threading import Event, Lock, Thread import pytz -import telnetlib3 from core.config import SERVER_OWNER_CALLSIGN from data.spot import Spot -from providers.spot.spot_provider import SpotProvider -from core.utils import decode_telnet_bytes +from providers.spot.telnet_spot_provider import TelnetSpotProvider logger = logging.getLogger(__name__) -class DXCluster(SpotProvider): +class DXCluster(TelnetSpotProvider): """Spot provider for a DX Cluster. Hostname, port, login_prompt, login_callsign and allow_rbn_spots are provided in config. See config-example.yml for examples.""" @@ -32,112 +28,33 @@ class DXCluster(SpotProvider): """Constructor requires hostname and port""" name = provider_config.get("name", "Cluster") - super().__init__(name, provider_config) - self._hostname = provider_config["host"] - self._port = provider_config["port"] - self._login_prompt = provider_config.get("login_prompt", "login:") - self._login_callsign = provider_config.get("login_callsign", SERVER_OWNER_CALLSIGN) - self._allow_rbn_spots = provider_config.get("allow_rbn_spots", False) - self._spot_line_pattern = ( - self._LINE_PATTERN_ALLOW_RBN if self._allow_rbn_spots else self._LINE_PATTERN_EXCLUDE_RBN + allow_rbn_spots = provider_config.get("allow_rbn_spots", False) + self._spot_line_pattern = self._LINE_PATTERN_ALLOW_RBN if allow_rbn_spots else self._LINE_PATTERN_EXCLUDE_RBN + super().__init__( + name, + provider_config, + host=provider_config["host"], + port=provider_config["port"], + login_prompt=provider_config.get("login_prompt", "login:"), + login_response=provider_config.get("login_callsign", SERVER_OWNER_CALLSIGN), ) - self._telnet = None - self._telnet_lock = Lock() - 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 _parse_line(self, line): + match = self._spot_line_pattern.match(line) + if not match: + return None - def stop(self): - self._stop_event.set() - with self._telnet_lock: - if self._telnet: - try: - self._telnet.sock.shutdown(socket.SHUT_RDWR) - except (AttributeError, OSError): - pass - self._telnet.close() - if self._thread: - self._thread.join(timeout=5) - 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 not self._stop_event.is_set(): - connected = False - while not connected and not self._stop_event.is_set(): - try: - self.status = "Connecting" - logger.info(f"DX Cluster {self._hostname} connecting...") - new_telnet = telnetlib3.Telnet(self._hostname, self._port) - with self._telnet_lock: - self._telnet = new_telnet - if self._stop_event.is_set(): - # stop() was called while we were connecting, close the connection rather than trying to - # read when we know it won't work - new_telnet.close() - break - self._telnet.read_until(self._login_prompt.encode("latin-1")) - self._telnet.write(f"{self._login_callsign}\n".encode("latin-1")) - connected = True - logger.info(f"DX Cluster {self._hostname} connected.") - except ConnectionRefusedError: - self.status = "Error" - logger.warning(f"Connection refused to DX cluster {self._hostname}") - self._stop_event.wait(timeout=300) - except Exception: - self.status = "Error" - logger.exception(f"Exception while connecting to DX Cluster Provider ({self._hostname}).") - self._stop_event.wait(timeout=5) - - self.status = "Waiting for Data" - 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")) - match = self._spot_line_pattern.match(decode_telnet_bytes(telnet_output)) - if match: - spot_time = datetime.strptime(match.group(5), "%H%MZ").replace(tzinfo=pytz.UTC) - spot_datetime = datetime.combine( - datetime.now(pytz.UTC).date(), - spot_time.time(), - tzinfo=pytz.UTC, - ) - spot = Spot( - source=self.name, - dx_call=match.group(3), - de_call=match.group(1), - freq=float(match.group(2)) * 1000, - comment=match.group(4).strip(), - time=spot_datetime.timestamp(), - ) - - # Add to our list - self._submit(spot) - - self.status = "OK" - self.last_update_time = datetime.now(pytz.UTC) - logger.debug(f"Data received from DX Cluster {self._hostname}.") - - except EOFError: - connected = False - if not self._stop_event.is_set(): - self.status = "Restarting" - logger.warning(f"Disconnected from DX Cluster {self._hostname}. Reconnecting...") - 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 not self._stop_event.is_set(): - self.status = "Error" - logger.exception(f"Exception in DX Cluster Provider ({self._hostname})") - self._stop_event.wait(timeout=5) - else: - logger.info(f"DX Cluster {self._hostname} shutting down...") - self.status = "Shutting down" - - self.status = "Disconnected" + spot_time = datetime.strptime(match.group(5), "%H%MZ").replace(tzinfo=pytz.UTC) + spot_datetime = datetime.combine( + datetime.now(pytz.UTC).date(), + spot_time.time(), + tzinfo=pytz.UTC, + ) + return Spot( + source=self.name, + dx_call=match.group(3), + de_call=match.group(1), + freq=float(match.group(2)) * 1000, + comment=match.group(4).strip(), + time=spot_datetime.timestamp(), + ) diff --git a/providers/spot/rbn.py b/providers/spot/rbn.py index bdfd78e..5862bcc 100644 --- a/providers/spot/rbn.py +++ b/providers/spot/rbn.py @@ -1,23 +1,19 @@ import logging import re -import socket from datetime import datetime -from threading import Event, Lock, Thread import pytz -import telnetlib3 from core.config import SERVER_OWNER_CALLSIGN from data.spot import Spot -from providers.spot.spot_provider import SpotProvider -from core.utils import decode_telnet_bytes +from providers.spot.telnet_spot_provider import TelnetSpotProvider logger = logging.getLogger(__name__) -class RBN(SpotProvider): +class RBN(TelnetSpotProvider): """Spot provider for the Reverse Beacon Network. Connects to a single port, if you want both CW/RTTY (port 7000) and FT8 - (port 7001) you need to instantiate two copies of this. The port is provided as an argument to the constructor.""" + (port 7001) you need to instantiate two copies of this. The port is provided in config.""" _LINE_PATTERN = re.compile( r"^DX de ([a-z0-9/]+)-.*:\s+([0-9.]+)\s+([a-z0-9/]+)\s+(.*)\s+(\d{4}Z)", @@ -28,101 +24,31 @@ class RBN(SpotProvider): """Constructor requires port number.""" name = provider_config.get("name", "RBN") - super().__init__(name, provider_config) - self._port = provider_config["port"] - self._telnet = None - self._telnet_lock = Lock() - self._thread = None - self._stop_event = Event() + super().__init__( + name, + provider_config, + host="telnet.reversebeacon.net", + port=provider_config["port"], + login_prompt="Please enter your call: ", + login_response=SERVER_OWNER_CALLSIGN, + ) - def start(self): - self._thread = Thread(target=self._handle, name=f"RBNSpotProvider-{self.name}", daemon=True) - self._thread.start() + def _parse_line(self, line): + match = self._LINE_PATTERN.match(line) + if not match: + return None - def stop(self): - self._stop_event.set() - with self._telnet_lock: - if self._telnet: - try: - self._telnet.sock.shutdown(socket.SHUT_RDWR) - except (AttributeError, OSError): - pass - self._telnet.close() - if self._thread: - self._thread.join(timeout=5) - 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 not self._stop_event.is_set(): - connected = False - while not connected and not self._stop_event.is_set(): - try: - self.status = "Connecting" - logger.info(f"RBN port {self._port!s} connecting...") - new_telnet = telnetlib3.Telnet("telnet.reversebeacon.net", self._port) - with self._telnet_lock: - self._telnet = new_telnet - if self._stop_event.is_set(): - # stop() was called while we were connecting, close the connection rather than trying to - # read when we know it won't work - new_telnet.close() - break - self._telnet.read_until("Please enter your call: ".encode("latin-1")) - self._telnet.write(f"{SERVER_OWNER_CALLSIGN}\n".encode("latin-1")) - connected = True - logger.info(f"RBN port {self._port!s} connected.") - except Exception: - self.status = "Error" - logger.exception(f"Exception while connecting to RBN (port {self._port!s}).") - self._stop_event.wait(timeout=5) - - self.status = "Waiting for Data" - 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")) - match = self._LINE_PATTERN.match(decode_telnet_bytes(telnet_output)) - if match: - spot_time = datetime.strptime(match.group(5), "%H%MZ").replace(tzinfo=pytz.UTC) - spot_datetime = datetime.combine( - datetime.now(pytz.UTC).date(), - spot_time.time(), - tzinfo=pytz.UTC, - ) - spot = Spot( - source=self.name, - dx_call=match.group(3), - de_call=match.group(1), - freq=float(match.group(2)) * 1000, - comment=match.group(4).strip(), - time=spot_datetime.timestamp(), - ) - - # Add to our list - self._submit(spot) - - self.status = "OK" - self.last_update_time = datetime.now(pytz.UTC) - logger.debug(f"Data received from RBN on port {self._port!s}.") - - except EOFError: - connected = False - if not self._stop_event.is_set(): - self.status = "Restarting" - logger.warning(f"Disconnected from RBN provider (port {self._port!s}). Reconnecting...") - 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 not self._stop_event.is_set(): - self.status = "Error" - logger.exception(f"Exception in RBN provider (port {self._port!s})") - self._stop_event.wait(timeout=5) - else: - logger.info(f"RBN provider (port {self._port!s}) shutting down...") - self.status = "Shutting down" - - self.status = "Disconnected" + spot_time = datetime.strptime(match.group(5), "%H%MZ").replace(tzinfo=pytz.UTC) + spot_datetime = datetime.combine( + datetime.now(pytz.UTC).date(), + spot_time.time(), + tzinfo=pytz.UTC, + ) + return Spot( + source=self.name, + dx_call=match.group(3), + de_call=match.group(1), + freq=float(match.group(2)) * 1000, + comment=match.group(4).strip(), + time=spot_datetime.timestamp(), + ) diff --git a/providers/spot/telnet_spot_provider.py b/providers/spot/telnet_spot_provider.py new file mode 100644 index 0000000..2c9e207 --- /dev/null +++ b/providers/spot/telnet_spot_provider.py @@ -0,0 +1,117 @@ +import logging +import socket +from datetime import datetime +from threading import Event, Lock, Thread + +import pytz +import telnetlib3 + +from core.utils import decode_telnet_bytes +from providers.spot.spot_provider import SpotProvider + +logger = logging.getLogger(__name__) + + +class TelnetSpotProvider(SpotProvider): + """Base class for spot providers that connect to a telnet server to receive spots.""" + + def __init__(self, name, provider_config, host, port, login_prompt, login_response): + """Constructor. host/port are the telnet server to connect to. login_prompt is the text to wait for before + logging in, and login_response is the callsign to send in response.""" + + super().__init__(name, provider_config) + self._host = host + self._port = port + self._login_prompt = login_prompt + self._login_response = login_response + self._telnet = None + self._telnet_lock = Lock() + self._thread = None + self._stop_event = Event() + + def _parse_line(self, line): + """Parse a line of telnet output and return a Spot, or None if the line did not contain a spot. Subclasses + must implement this method.""" + + raise NotImplementedError("Subclasses must implement this method") + + def start(self): + self._thread = Thread(target=self._handle, name=f"{self.name}SpotProvider", daemon=True) + self._thread.start() + + def stop(self): + self._stop_event.set() + with self._telnet_lock: + if self._telnet: + try: + self._telnet.sock.shutdown(socket.SHUT_RDWR) + except (AttributeError, OSError): + pass + self._telnet.close() + if self._thread: + self._thread.join(timeout=5) + if self._thread.is_alive(): + logger.warning(f"{self.name} worker thread did not exit on time and will be killed.") + + def _handle(self): + while not self._stop_event.is_set(): + connected = False + while not connected and not self._stop_event.is_set(): + try: + self.status = "Connecting" + logger.info(f"{self.name} connecting to {self._host}:{self._port!s}...") + new_telnet = telnetlib3.Telnet(self._host, self._port) + with self._telnet_lock: + self._telnet = new_telnet + if self._stop_event.is_set(): + # stop() was called while we were connecting, close the connection rather than trying to + # read when we know it won't work + new_telnet.close() + break + self._telnet.read_until(self._login_prompt.encode("latin-1")) + self._telnet.write(f"{self._login_response}\n".encode("latin-1")) + connected = True + logger.info(f"{self.name} connected.") + except ConnectionRefusedError: + self.status = "Error" + logger.warning(f"Connection refused to {self.name} ({self._host}:{self._port!s}).") + self._stop_event.wait(timeout=300) + except Exception: + self.status = "Error" + logger.exception(f"Exception while connecting to {self.name} ({self._host}:{self._port!s}).") + self._stop_event.wait(timeout=5) + + self.status = "Waiting for Data" + 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")) + spot = self._parse_line(decode_telnet_bytes(telnet_output)) + if spot: + # Add to our list + self._submit(spot) + + self.status = "OK" + self.last_update_time = datetime.now(pytz.UTC) + logger.debug(f"Data received from {self.name}.") + + except EOFError: + connected = False + if not self._stop_event.is_set(): + self.status = "Restarting" + logger.warning(f"Disconnected from {self.name}. Reconnecting...") + self._stop_event.wait(timeout=5) + else: + logger.info(f"{self.name} shutting down...") + self.status = "Shutting down" + except Exception: + connected = False + if not self._stop_event.is_set(): + self.status = "Error" + logger.exception(f"Exception in {self.name}") + self._stop_event.wait(timeout=5) + else: + logger.info(f"{self.name} shutting down...") + self.status = "Shutting down" + + self.status = "Disconnected" diff --git a/templates/add_spot.html b/templates/add_spot.html index 0ce2af2..294d7c7 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 2eb1a07..3009080 100644 --- a/templates/alerts.html +++ b/templates/alerts.html @@ -85,7 +85,7 @@ - + diff --git a/templates/bands.html b/templates/bands.html index 42b68fe..65e5ec0 100644 --- a/templates/bands.html +++ b/templates/bands.html @@ -76,8 +76,8 @@ - - + + diff --git a/templates/base.html b/templates/base.html index f302c60..42b53cb 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 %}