diff --git a/config.template b/config.template index c705252..1c64f72 100644 --- a/config.template +++ b/config.template @@ -335,6 +335,20 @@ voxOnTrapList = True voxTrapList = chirpy voxEnableCmd = True +# WSJT-X UDP monitoring - listens for decode messages from WSJT-X, FT8/FT4/WSPR etc. +wsjtxDetectionEnabled = False +# UDP address and port where WSJT-X broadcasts (default: 127.0.0.1:2237) +wsjtxUdpServerAddress = 127.0.0.1:2237 +# Comma-separated list of callsigns to watch (empty = all callsigns) +wsjtxWatchedCallsigns = + +# JS8Call TCP monitoring - connects to JS8Call API for message forwarding +js8callDetectionEnabled = False +# TCP address and port where JS8Call API listens (default: 127.0.0.1:2442) +js8callServerAddress = 127.0.0.1:2442 +# Comma-separated list of callsigns to watch (empty = all callsigns) +js8callWatchedCallsigns = + [fileMon] filemon_enabled = False diff --git a/mesh_bot.py b/mesh_bot.py index 85becdf..c9a5229 100755 --- a/mesh_bot.py +++ b/mesh_bot.py @@ -2020,6 +2020,12 @@ async def main(): if my_settings.voxDetectionEnabled: tasks.append(asyncio.create_task(voxMonitor(), name="vox_detection")) + + if my_settings.wsjtx_detection_enabled: + tasks.append(asyncio.create_task(handleWsjtxWatcher(), name="wsjtx_monitor")) + + if my_settings.js8call_detection_enabled: + tasks.append(asyncio.create_task(handleJs8callWatcher(), name="js8call_monitor")) if my_settings.scheduler_enabled: from modules.scheduler import run_scheduler_loop, setup_scheduler diff --git a/modules/radio.py b/modules/radio.py index 4b1cb28..89350b6 100644 --- a/modules/radio.py +++ b/modules/radio.py @@ -263,4 +263,309 @@ async def voxMonitor(): except Exception as e: logger.error(f"RadioMon: Error in VOX monitor: {e}") +# WSJT-X and JS8Call UDP Monitoring +# Based on WSJT-X UDP protocol specification +# Reference: https://github.com/ckuhtz/ham/blob/main/mcast/recv_decode.py + +wsjtx_enabled = False +js8call_enabled = False +wsjtx_udp_port = 2237 +js8call_udp_port = 2442 +watched_callsigns = [] +wsjtx_udp_address = '127.0.0.1' +js8call_tcp_address = '127.0.0.1' +js8call_tcp_port = 2442 + +try: + from modules.settings import ( + wsjtx_detection_enabled, + wsjtx_udp_server_address, + wsjtx_watched_callsigns, + js8call_detection_enabled, + js8call_server_address, + js8call_watched_callsigns + ) + wsjtx_enabled = wsjtx_detection_enabled + js8call_enabled = js8call_detection_enabled + + if wsjtx_enabled: + import socket + import struct + # Parse UDP address + if ':' in wsjtx_udp_server_address: + wsjtx_udp_address, port_str = wsjtx_udp_server_address.split(':') + wsjtx_udp_port = int(port_str) + watched_callsigns.extend(wsjtx_watched_callsigns.split(',') if wsjtx_watched_callsigns else []) + + if js8call_enabled: + import socket + import json + # Parse TCP address for JS8Call + if ':' in js8call_server_address: + js8call_tcp_address, port_str = js8call_server_address.split(':') + js8call_tcp_port = int(port_str) + watched_callsigns.extend(js8call_watched_callsigns.split(',') if js8call_watched_callsigns else []) + + # Clean up callsigns - remove whitespace + watched_callsigns = [cs.strip().upper() for cs in watched_callsigns if cs.strip()] + +except ImportError: + logger.debug("RadioMon: WSJT-X/JS8Call settings not configured") +except Exception as e: + logger.warning(f"RadioMon: Error loading WSJT-X/JS8Call settings: {e}") + +# WSJT-X UDP Protocol Message Types +WSJTX_HEARTBEAT = 0 +WSJTX_STATUS = 1 +WSJTX_DECODE = 2 +WSJTX_CLEAR = 3 +WSJTX_REPLY = 4 +WSJTX_QSO_LOGGED = 5 +WSJTX_CLOSE = 6 +WSJTX_REPLAY = 7 +WSJTX_HALT_TX = 8 +WSJTX_FREE_TEXT = 9 +WSJTX_WSPR_DECODE = 10 +WSJTX_LOCATION = 11 +WSJTX_LOGGED_ADIF = 12 + +wsjtxMsgQueue = [] # Queue for WSJT-X detected messages +js8callMsgQueue = [] # Queue for JS8Call detected messages + +def decode_wsjtx_packet(data): + """Decode WSJT-X UDP packet according to the protocol specification""" + try: + # WSJT-X uses Qt's QDataStream format (big-endian) + magic = struct.unpack('>I', data[0:4])[0] + if magic != 0xADBCCBDA: + return None + + schema_version = struct.unpack('>I', data[4:8])[0] + msg_type = struct.unpack('>I', data[8:12])[0] + + offset = 12 + + # Helper to read Qt QString (4-byte length + UTF-8 data) + def read_qstring(data, offset): + if offset + 4 > len(data): + return "", offset + length = struct.unpack('>I', data[offset:offset+4])[0] + offset += 4 + if length == 0xFFFFFFFF: # Null string + return "", offset + if offset + length > len(data): + return "", offset + text = data[offset:offset+length].decode('utf-8', errors='ignore') + return text, offset + length + + # Decode DECODE message (type 2) + if msg_type == WSJTX_DECODE: + # Read fields according to WSJT-X protocol + wsjtx_id, offset = read_qstring(data, offset) + + # Read other decode fields: new, time, snr, delta_time, delta_frequency, mode, message + if offset + 1 > len(data): + return None + new = struct.unpack('>?', data[offset:offset+1])[0] + offset += 1 + + if offset + 4 > len(data): + return None + time_val = struct.unpack('>I', data[offset:offset+4])[0] + offset += 4 + + if offset + 4 > len(data): + return None + snr = struct.unpack('>i', data[offset:offset+4])[0] + offset += 4 + + if offset + 8 > len(data): + return None + delta_time = struct.unpack('>d', data[offset:offset+8])[0] + offset += 8 + + if offset + 4 > len(data): + return None + delta_frequency = struct.unpack('>I', data[offset:offset+4])[0] + offset += 4 + + mode, offset = read_qstring(data, offset) + message, offset = read_qstring(data, offset) + + return { + 'type': 'decode', + 'id': wsjtx_id, + 'new': new, + 'time': time_val, + 'snr': snr, + 'delta_time': delta_time, + 'delta_frequency': delta_frequency, + 'mode': mode, + 'message': message + } + + # Decode QSO_LOGGED message (type 5) + elif msg_type == WSJTX_QSO_LOGGED: + wsjtx_id, offset = read_qstring(data, offset) + + # Read QSO logged fields + if offset + 8 > len(data): + return None + date_off = struct.unpack('>Q', data[offset:offset+8])[0] + offset += 8 + + if offset + 8 > len(data): + return None + time_off = struct.unpack('>Q', data[offset:offset+8])[0] + offset += 8 + + dx_call, offset = read_qstring(data, offset) + dx_grid, offset = read_qstring(data, offset) + + return { + 'type': 'qso_logged', + 'id': wsjtx_id, + 'dx_call': dx_call, + 'dx_grid': dx_grid + } + + return None + + except Exception as e: + logger.debug(f"RadioMon: Error decoding WSJT-X packet: {e}") + return None + +def check_callsign_match(message, callsigns): + """Check if any watched callsign appears in the message""" + if not callsigns: + return True # If no filter, accept all + + message_upper = message.upper() + for callsign in callsigns: + if callsign in message_upper: + return True + return False + +async def wsjtxMonitor(): + """Monitor WSJT-X UDP broadcasts for decode messages""" + if not wsjtx_enabled: + logger.warning("RadioMon: WSJT-X monitoring called but not enabled") + return + + try: + sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) + sock.bind((wsjtx_udp_address, wsjtx_udp_port)) + sock.setblocking(False) + + logger.info(f"RadioMon: WSJT-X UDP listener started on {wsjtx_udp_address}:{wsjtx_udp_port}") + if watched_callsigns: + logger.info(f"RadioMon: Watching for callsigns: {', '.join(watched_callsigns)}") + + while True: + try: + data, addr = sock.recvfrom(4096) + decoded = decode_wsjtx_packet(data) + + if decoded and decoded['type'] == 'decode': + message = decoded['message'] + mode = decoded['mode'] + snr = decoded['snr'] + + # Check if message contains watched callsigns + if check_callsign_match(message, watched_callsigns): + msg_text = f"WSJT-X {mode}: {message} (SNR: {snr:+d}dB)" + logger.info(f"RadioMon: {msg_text}") + wsjtxMsgQueue.append(msg_text) + + except BlockingIOError: + # No data available + await asyncio.sleep(0.1) + except Exception as e: + logger.debug(f"RadioMon: Error in WSJT-X monitor loop: {e}") + await asyncio.sleep(1) + + except Exception as e: + logger.error(f"RadioMon: Error starting WSJT-X monitor: {e}") + +async def js8callMonitor(): + """Monitor JS8Call TCP API for messages""" + if not js8call_enabled: + logger.warning("RadioMon: JS8Call monitoring called but not enabled") + return + + try: + logger.info(f"RadioMon: JS8Call TCP listener connecting to {js8call_tcp_address}:{js8call_tcp_port}") + if watched_callsigns: + logger.info(f"RadioMon: Watching for callsigns: {', '.join(watched_callsigns)}") + + while True: + try: + # Connect to JS8Call TCP API + sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) + sock.settimeout(5) + sock.connect((js8call_tcp_address, js8call_tcp_port)) + sock.setblocking(False) + + logger.info("RadioMon: Connected to JS8Call API") + + buffer = "" + while True: + try: + data = sock.recv(4096) + if not data: + logger.warning("RadioMon: JS8Call connection closed") + break + + buffer += data.decode('utf-8', errors='ignore') + + # Process complete JSON messages (newline delimited) + while '\n' in buffer: + line, buffer = buffer.split('\n', 1) + if not line.strip(): + continue + + try: + msg = json.loads(line) + msg_type = msg.get('type', '') + + # Handle RX.DIRECTED and RX.ACTIVITY messages + if msg_type in ['RX.DIRECTED', 'RX.ACTIVITY']: + params = msg.get('params', {}) + text = params.get('TEXT', '') + from_call = params.get('FROM', '') + snr = params.get('SNR', 0) + + if text and check_callsign_match(text, watched_callsigns): + msg_text = f"JS8Call from {from_call}: {text} (SNR: {snr:+d}dB)" + logger.info(f"RadioMon: {msg_text}") + js8callMsgQueue.append(msg_text) + + except json.JSONDecodeError: + logger.debug(f"RadioMon: Invalid JSON from JS8Call: {line[:100]}") + except Exception as e: + logger.debug(f"RadioMon: Error processing JS8Call message: {e}") + + except BlockingIOError: + await asyncio.sleep(0.1) + except socket.timeout: + await asyncio.sleep(0.1) + except Exception as e: + logger.debug(f"RadioMon: Error in JS8Call receive loop: {e}") + break + + sock.close() + logger.warning("RadioMon: JS8Call connection lost, reconnecting in 5s...") + await asyncio.sleep(5) + + except socket.timeout: + logger.warning("RadioMon: JS8Call connection timeout, retrying in 5s...") + await asyncio.sleep(5) + except Exception as e: + logger.warning(f"RadioMon: Error connecting to JS8Call: {e}") + await asyncio.sleep(10) + + except Exception as e: + logger.error(f"RadioMon: Error starting JS8Call monitor: {e}") + # end of file diff --git a/modules/settings.py b/modules/settings.py index ae89ae1..5cdc327 100644 --- a/modules/settings.py +++ b/modules/settings.py @@ -32,6 +32,8 @@ cmdHistory = [] # list to hold the command history for lheard and history comman msg_history = [] # list to hold the message history for the messages command max_bytes = 200 # Meshtastic has ~237 byte limit, use conservative 200 bytes for message content voxMsgQueue = [] # queue for VOX detected messages +wsjtxMsgQueue = [] # queue for WSJT-X detected messages +js8callMsgQueue = [] # queue for JS8Call detected messages # Game trackers surveyTracker = [] # Survey game tracker tictactoeTracker = [] # TicTacToe game tracker @@ -406,6 +408,14 @@ try: voxOnTrapList = config['radioMon'].getboolean('voxOnTrapList', False) # default False voxTrapList = config['radioMon'].get('voxTrapList', 'chirpy').split(',') # default chirpy voxEnableCmd = config['radioMon'].getboolean('voxEnableCmd', True) # default True + + # WSJT-X and JS8Call monitoring + wsjtx_detection_enabled = config['radioMon'].getboolean('wsjtxDetectionEnabled', False) # default WSJT-X detection disabled + wsjtx_udp_server_address = config['radioMon'].get('wsjtxUdpServerAddress', '127.0.0.1:2237') # default localhost:2237 + wsjtx_watched_callsigns = config['radioMon'].get('wsjtxWatchedCallsigns', '') # default empty (all callsigns) + js8call_detection_enabled = config['radioMon'].getboolean('js8callDetectionEnabled', False) # default JS8Call detection disabled + js8call_server_address = config['radioMon'].get('js8callServerAddress', '127.0.0.1:2442') # default localhost:2442 + js8call_watched_callsigns = config['radioMon'].get('js8callWatchedCallsigns', '') # default empty (all callsigns) # file monitor file_monitor_enabled = config['fileMon'].getboolean('filemon_enabled', False) diff --git a/modules/system.py b/modules/system.py index 9b0eb5b..3373263 100644 --- a/modules/system.py +++ b/modules/system.py @@ -2007,6 +2007,62 @@ async def handleFileWatcher(): await asyncio.sleep(1) pass +async def handleWsjtxWatcher(): + # monitor WSJT-X UDP broadcasts for decode messages + from modules.radio import wsjtxMsgQueue, wsjtxMonitor + from modules.settings import sigWatchBroadcastCh, sigWatchBroadcastInterface + + # Start the WSJT-X monitor task + monitor_task = asyncio.create_task(wsjtxMonitor()) + + while True: + if wsjtxMsgQueue: + msg = wsjtxMsgQueue.pop(0) + logger.debug(f"System: Detected message from WSJT-X: {msg}") + + # Broadcast to configured channels + if type(sigWatchBroadcastCh) is list: + for ch in sigWatchBroadcastCh: + if antiSpam and int(ch) != publicChannel: + send_message(msg, int(ch), 0, sigWatchBroadcastInterface) + else: + logger.warning(f"System: antiSpam prevented Alert from WSJT-X") + else: + if antiSpam and sigWatchBroadcastCh != publicChannel: + send_message(msg, int(sigWatchBroadcastCh), 0, sigWatchBroadcastInterface) + else: + logger.warning(f"System: antiSpam prevented Alert from WSJT-X") + + await asyncio.sleep(0.5) + +async def handleJs8callWatcher(): + # monitor JS8Call TCP API for messages + from modules.radio import js8callMsgQueue, js8callMonitor + from modules.settings import sigWatchBroadcastCh, sigWatchBroadcastInterface + + # Start the JS8Call monitor task + monitor_task = asyncio.create_task(js8callMonitor()) + + while True: + if js8callMsgQueue: + msg = js8callMsgQueue.pop(0) + logger.debug(f"System: Detected message from JS8Call: {msg}") + + # Broadcast to configured channels + if type(sigWatchBroadcastCh) is list: + for ch in sigWatchBroadcastCh: + if antiSpam and int(ch) != publicChannel: + send_message(msg, int(ch), 0, sigWatchBroadcastInterface) + else: + logger.warning(f"System: antiSpam prevented Alert from JS8Call") + else: + if antiSpam and sigWatchBroadcastCh != publicChannel: + send_message(msg, int(sigWatchBroadcastCh), 0, sigWatchBroadcastInterface) + else: + logger.warning(f"System: antiSpam prevented Alert from JS8Call") + + await asyncio.sleep(0.5) + async def retry_interface(nodeID): global retry_int1, retry_int2, retry_int3, retry_int4, retry_int5, retry_int6, retry_int7, retry_int8, retry_int9 global max_retry_count1, max_retry_count2, max_retry_count3, max_retry_count4, max_retry_count5, max_retry_count6, max_retry_count7, max_retry_count8, max_retry_count9