Extract some common code into telnet_spot_provier as a base class for DX cluster and RBN

This commit is contained in:
Ian Renton
2026-09-22 18:35:05 +01:00
parent 9067f545af
commit f8343dd036
11 changed files with 190 additions and 230 deletions
+29 -112
View File
@@ -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(),
)
+29 -103
View File
@@ -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(),
)
+117
View File
@@ -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"