Implement text message handling and identity registration for room servers

This commit is contained in:
Lloyd
2025-12-17 10:38:23 +00:00
parent 73d5250560
commit 4b97c7b439
7 changed files with 326 additions and 63 deletions
+7 -2
View File
@@ -63,7 +63,6 @@ mesh:
# Multiple Identity Configuration (Optional)
# Define additional identities for the repeater to manage
# Each identity operates independently with its own key pair and configuration
# Note: All identities share the same ACL (security settings from repeater.security)
identities:
# Room Server Identities
# Each room server acts as a separate logical node on the mesh
@@ -78,7 +77,10 @@ identities:
# node_name: "Test BBS Room"
# latitude: 0.0
# longitude: 0.0
# disable_fwd: true # Room servers typically don't forward
#
# security:
# admin_password: "room_admin_password"
# guest_password: "room_guest_password"
# Add more room servers as needed
# - name: "SocialHub"
@@ -86,6 +88,9 @@ identities:
# type: "room_server"
# settings:
# node_name: "Social Hub"
# security:
# admin_password: "social_admin_123"
# guest_password: "social_guest_123"
radio:
# Frequency in Hz (869.618 MHz for EU)
+2 -1
View File
@@ -4,5 +4,6 @@ from .trace import TraceHelper
from .discovery import DiscoveryHelper
from .advert import AdvertHelper
from .login import LoginHelper
from .text import TextHelper
__all__ = ["TraceHelper", "DiscoveryHelper", "AdvertHelper", "LoginHelper"]
__all__ = ["TraceHelper", "DiscoveryHelper", "AdvertHelper", "LoginHelper", "TextHelper"]
+38 -13
View File
@@ -1,10 +1,3 @@
"""
Access Control List for pyMC Repeater.
Manages authenticated clients with permission-based access control.
Shared across all identities (repeater and room servers).
"""
import logging
import time
from typing import Dict, Optional
@@ -56,8 +49,40 @@ class ACL:
self.clients: Dict[bytes, ClientInfo] = {}
def authenticate_client(
self, client_identity: Identity, shared_secret: bytes, password: str, timestamp: int
self,
client_identity: Identity,
shared_secret: bytes,
password: str,
timestamp: int,
target_identity_hash: int = None,
target_identity_name: str = None,
target_identity_config: dict = None
) -> tuple[bool, int]:
target_identity_config = target_identity_config or {}
# Check for identity-specific passwords (required for room servers)
identity_security = target_identity_config.get("security", {})
# Determine if this is a room server by checking the type field
identity_type = target_identity_config.get("type", "")
is_room_server = identity_type == "room_server"
if is_room_server:
# Room servers MUST define their own passwords - no fallback
admin_pwd = identity_security.get("admin_password")
guest_pwd = identity_security.get("guest_password")
if not admin_pwd and not guest_pwd:
logger.error(f"Room server '{target_identity_name}' has no passwords configured!")
return False, 0
else:
# Repeater uses global passwords
admin_pwd = self.admin_password
guest_pwd = self.guest_password
if target_identity_name:
logger.debug(f"Authenticating for identity '{target_identity_name}' (room_server={is_room_server})")
pub_key = client_identity.get_public_key()[:PUB_KEY_SIZE]
@@ -74,14 +99,14 @@ class ACL:
return True, client.permissions
permissions = 0
if password == self.admin_password:
if admin_pwd and password == admin_pwd:
permissions = PERM_ACL_ADMIN
logger.info("Admin password validated")
elif self.guest_password and password == self.guest_password:
logger.info(f"Admin password validated for '{target_identity_name or 'unknown'}'")
elif guest_pwd and password == guest_pwd:
permissions = PERM_ACL_READ_WRITE
logger.info("Guest password validated")
logger.info(f"Guest password validated for '{target_identity_name or 'unknown'}'")
else:
logger.info("Invalid password")
logger.info(f"Invalid password for '{target_identity_name or 'unknown'}'")
return False, 0
client = self.clients.get(pub_key)
+30 -6
View File
@@ -13,8 +13,6 @@ logger = logging.getLogger("LoginHelper")
class LoginHelper:
"""Helper class for processing ANON_REQ login packets in the repeater."""
def __init__(self, identity_manager, packet_injector=None, acl=None, log_fn=None):
self.identity_manager = identity_manager
@@ -24,25 +22,51 @@ class LoginHelper:
self.handlers = {}
def register_identity(self, name: str, identity, identity_type: str = "room_server"):
def register_identity(self, name: str, identity, identity_type: str = "room_server", config: dict = None):
if not self.acl:
logger.warning(f"Cannot register identity '{name}': no ACL configured")
return
config = config or {}
# Validate room servers have their own passwords
if identity_type == "room_server":
security = config.get("security", {})
if not security.get("admin_password") and not security.get("guest_password"):
logger.error(
f"Room server '{name}' MUST have security.admin_password or "
f"security.guest_password configured. Skipping registration."
)
return
hash_byte = identity.get_public_key()[0]
# Create auth callback that includes identity context
def auth_callback_with_context(client_identity, shared_secret, password, timestamp):
return self.acl.authenticate_client(
client_identity=client_identity,
shared_secret=shared_secret,
password=password,
timestamp=timestamp,
target_identity_hash=hash_byte, # Which identity is being logged into
target_identity_name=name,
target_identity_config=config
)
handler = LoginServerHandler(
local_identity=identity,
log_fn=self.log_fn,
authenticate_callback=self.acl.authenticate_client,
authenticate_callback=auth_callback_with_context,
)
handler.set_send_packet_callback(self._send_packet_with_delay)
hash_byte = identity.get_public_key()[0]
self.handlers[hash_byte] = handler
logger.info(f"Registered {identity_type} '{name}' login handler: hash=0x{hash_byte:02X}")
async def process_login_packet(self, packet):
try:
+174
View File
@@ -0,0 +1,174 @@
"""
Text message (TXT_MSG) handling helper for pyMC Repeater.
This module processes incoming text messages for all managed identities
(repeater identity + identity manager identities).
"""
import asyncio
import logging
from pymc_core.node.handlers.text import TextMessageHandler
logger = logging.getLogger("TextHelper")
class TextHelper:
def __init__(self, identity_manager, packet_injector=None, acl=None, log_fn=None):
self.identity_manager = identity_manager
self.packet_injector = packet_injector
self.log_fn = log_fn or logger.info
self.acl = acl
# Dictionary of handlers keyed by dest_hash
self.handlers = {}
def register_identity(
self,
name: str,
identity,
identity_type: str = "room_server",
radio_config=None
):
if not self.acl:
logger.warning(f"Cannot register identity '{name}': no ACL configured")
return
# Create a contacts wrapper from ACL
acl_contacts = self._create_acl_contacts_wrapper()
# Create TextMessageHandler for this identity
handler = TextMessageHandler(
local_identity=identity,
contacts=acl_contacts,
log_fn=self.log_fn,
send_packet_fn=self._send_packet,
radio_config=radio_config,
)
# Register by dest hash
hash_byte = identity.get_public_key()[0]
self.handlers[hash_byte] = {
"handler": handler,
"identity": identity,
"name": name,
"type": identity_type,
}
logger.info(
f"Registered {identity_type} '{name}' text handler: hash=0x{hash_byte:02X}"
)
def _create_acl_contacts_wrapper(self):
class ACLContactsWrapper:
def __init__(self, acl):
self._acl = acl
@property
def contacts(self):
contact_list = []
for client_info in self._acl.get_all_clients():
# Create a minimal contact object that TextMessageHandler needs
class ContactProxy:
def __init__(self, client):
self.public_key = client.id.get_public_key().hex()
self.name = f"client_{self.public_key[:8]}"
contact_list.append(ContactProxy(client_info))
return contact_list
return ACLContactsWrapper(self.acl)
async def process_text_packet(self, packet):
try:
if len(packet.payload) < 2:
return False
dest_hash = packet.payload[0]
src_hash = packet.payload[1]
handler_info = self.handlers.get(dest_hash)
if handler_info:
logger.debug(
f"Routing text message to '{handler_info['name']}': "
f"dest=0x{dest_hash:02X}, src=0x{src_hash:02X}"
)
# Call the handler
await handler_info["handler"](packet)
# Call placeholder for custom processing
await self._on_message_received(
identity_name=handler_info["name"],
identity_type=handler_info["type"],
packet=packet,
dest_hash=dest_hash,
src_hash=src_hash,
)
# Mark packet as handled
packet.mark_do_not_retransmit()
return True
else:
logger.debug(
f"No text handler for hash 0x{dest_hash:02X}, allowing forward"
)
return False
except Exception as e:
logger.error(f"Error processing text packet: {e}")
return False
async def _on_message_received(
self,
identity_name: str,
identity_type: str,
packet,
dest_hash: int,
src_hash: int,
):
# Placeholder - can be overridden or callback can be added
logger.debug(
f"Message received for {identity_type} '{identity_name}' "
f"from 0x{src_hash:02X}"
)
# Example: Extract decrypted message if available
if hasattr(packet, "decrypted") and packet.decrypted:
message_text = packet.decrypted.get("text", "<unknown>")
logger.info(
f"[{identity_type}:{identity_name}] Message: {message_text}"
)
async def _send_packet(self, packet, wait_for_ack: bool = False):
if self.packet_injector:
try:
return await self.packet_injector(packet, wait_for_ack=wait_for_ack)
except Exception as e:
logger.error(f"Error sending packet: {e}")
return False
else:
logger.error("No packet injector configured, cannot send packet")
return False
def set_message_callback(self, callback):
self._message_callback = callback
def list_registered_identities(self):
return [
{
"hash": hash_byte,
"name": info["name"],
"type": info["type"],
}
for hash_byte, info in self.handlers.items()
]
+64 -4
View File
@@ -6,7 +6,7 @@ import sys
from repeater.config import get_radio_for_board, load_config
from repeater.engine import RepeaterHandler
from repeater.web.http_server import HTTPStatsServer, _log_buffer
from repeater.handler_helpers import TraceHelper, DiscoveryHelper, AdvertHelper, LoginHelper
from repeater.handler_helpers import TraceHelper, DiscoveryHelper, AdvertHelper, LoginHelper, TextHelper
from repeater.packet_router import PacketRouter
from repeater.identity_manager import IdentityManager
@@ -29,6 +29,7 @@ class RepeaterDaemon:
self.advert_helper = None
self.discovery_helper = None
self.login_helper = None
self.text_helper = None
self.acl = None
self.router = None
@@ -184,14 +185,47 @@ class RepeaterDaemon:
self.login_helper.register_identity(
name="repeater",
identity=self.local_identity,
identity_type="repeater"
identity_type="repeater",
config=repeater_config # Pass repeater config (includes security settings)
)
# Register room server identities
# Register room server identities with their configs
for name, identity, config in self.identity_manager.get_identities_by_type("room_server"):
self.login_helper.register_identity(name, identity, identity_type="room_server")
self.login_helper.register_identity(
name=name,
identity=identity,
identity_type="room_server",
config=config # Pass room-specific config
)
logger.info("Login processing helper initialized")
# Initialize text message helper (uses shared ACL as contacts)
self.text_helper = TextHelper(
identity_manager=self.identity_manager,
packet_injector=self.router.inject_packet,
acl=shared_acl, # Use shared ACL as contacts database
log_fn=logger.info,
)
# Register default repeater identity for text messages
self.text_helper.register_identity(
name="repeater",
identity=self.local_identity,
identity_type="repeater",
radio_config=self.config.get("radio", {})
)
# Register room server identities for text messages
for name, identity, config in self.identity_manager.get_identities_by_type("room_server"):
self.text_helper.register_identity(
name=name,
identity=identity,
identity_type="room_server",
radio_config=self.config.get("radio", {})
)
logger.info("Text message processing helper initialized")
except Exception as e:
logger.error(f"Failed to initialize dispatcher: {e}")
@@ -250,6 +284,32 @@ class RepeaterDaemon:
await self.router.enqueue(packet)
except Exception as e:
logger.error(f"Error enqueuing packet in router: {e}", exc_info=True)
def register_text_handler_for_identity(
self,
name: str,
identity,
identity_type: str = "room_server",
radio_config: dict = None
):
if not self.text_helper:
logger.warning("Text helper not initialized, cannot register identity")
return False
try:
self.text_helper.register_identity(
name=name,
identity=identity,
identity_type=identity_type,
radio_config=radio_config or self.config.get("radio", {}),
)
logger.info(f"Registered text handler for {identity_type} '{name}'")
return True
except Exception as e:
logger.error(f"Failed to register text handler for '{name}': {e}")
return False
def get_stats(self) -> dict:
stats = {}
+11 -37
View File
@@ -1,11 +1,3 @@
"""
Packet router for pyMC Repeater.
This module provides a simple router that routes packets to appropriate handlers
based on payload type. All statistics, queuing, and processing logic is handled
by the repeater engine for better separation of concerns.
"""
import asyncio
import logging
@@ -13,16 +5,12 @@ from pymc_core.node.handlers.trace import TraceHandler
from pymc_core.node.handlers.control import ControlHandler
from pymc_core.node.handlers.advert import AdvertHandler
from pymc_core.node.handlers.login_server import LoginServerHandler
from pymc_core.node.handlers.text import TextMessageHandler
logger = logging.getLogger("PacketRouter")
class PacketRouter:
"""
Simple router that processes packets through handlers sequentially.
All statistics and processing decisions are handled by the engine.
"""
def __init__(self, daemon_instance):
self.daemon = daemon_instance
self.queue = asyncio.Queue()
@@ -30,13 +18,11 @@ class PacketRouter:
self.router_task = None
async def start(self):
"""Start the router processing task."""
self.running = True
self.router_task = asyncio.create_task(self._process_queue())
logger.info("Packet router started")
async def stop(self):
"""Stop the router processing task."""
self.running = False
if self.router_task:
self.router_task.cancel()
@@ -51,19 +37,6 @@ class PacketRouter:
await self.queue.put(packet)
async def inject_packet(self, packet, wait_for_ack: bool = False):
"""
Inject a new packet into the system for transmission through the engine.
This method uses the engine's main packet handler with the local_transmission
flag to bypass forwarding logic while maintaining proper statistics and airtime.
Args:
packet: The packet to send
wait_for_ack: Whether to wait for acknowledgment
Returns:
True if packet was sent successfully, False otherwise
"""
try:
metadata = {
"rssi": getattr(packet, "rssi", 0),
@@ -83,7 +56,6 @@ class PacketRouter:
return False
async def _process_queue(self):
"""Process packets through the router queue."""
while self.running:
try:
packet = await asyncio.wait_for(self.queue.get(), timeout=0.1)
@@ -95,13 +67,7 @@ class PacketRouter:
async def _route_packet(self, packet):
"""
Route a packet to appropriate handlers based on payload type.
Simple routing logic:
1. Route to specific handlers for parsing
2. Pass to repeater engine for all processing decisions
"""
payload_type = packet.get_payload_type()
processed_by_injection = False
@@ -134,6 +100,14 @@ class PacketRouter:
if handled:
processed_by_injection = True
elif payload_type == TextMessageHandler.payload_type():
# Process TXT_MSG packet for all identities
if self.daemon.text_helper:
handled = await self.daemon.text_helper.process_text_packet(packet)
# Only skip forwarding if we actually handled it
if handled:
processed_by_injection = True
# Only pass to repeater engine if not already processed by injection
if self.daemon.repeater_handler and not processed_by_injection:
metadata = {