feat(inject): realtime MySQL OPTIONS, dashboard push, and SINGLE=0 dynamic RX

Push STATE_SND when peers send RPTO so monitor chips update without reload.
Fetch and inject MySQL OPTIONS immediately on PASS= instead of a 10s delay.
Track per-peer dynamic TG sets for SINGLE=0 so keyed hotspots hear each other.
pull/1/head
Rodrigo Pérez 4 months ago
parent c59757e0ec
commit baaa0c3b00

@ -194,6 +194,61 @@ def peer_single_mode(peer: dict[str, Any], sys_cfg: dict[str, Any]) -> bool:
return single
def _peer_ua_multi_store(sys_cfg: dict[str, Any]) -> dict[bytes, dict[int, set[int]]]:
store = sys_cfg.setdefault("_PEER_UA_MULTI_TGS", {})
if not isinstance(store, dict):
store = {}
sys_cfg["_PEER_UA_MULTI_TGS"] = store
return store
def register_peer_ua_multi_tg(
peer: dict[str, Any],
peer_id: bytes,
slot: int,
tgid: int,
sys_cfg: dict[str, Any],
) -> None:
"""SINGLE=0: accumulate keyed dynamic TGs per peer/slot until TG 4000."""
if peer_single_mode(peer, sys_cfg):
return
tgid_i = int(tgid)
if tgid_i <= 0 or tgid_i == 4000:
return
if peer_receives_group_tgid(peer, slot, tgid_i):
return
pk = bytes_4(int_id(peer_id))
per_peer = _peer_ua_multi_store(sys_cfg).setdefault(pk, {})
slot_set = per_peer.setdefault(int(slot), set())
slot_set.add(tgid_i)
def peer_owns_multi_dynamic_ua(
peer: dict[str, Any],
slot: int,
tgid: int,
sys_cfg: dict[str, Any] | None,
*,
peer_id: bytes | None = None,
) -> bool:
"""True when SINGLE=0 peer has keyed this non-static dynamic TG on ``slot``."""
if not sys_cfg or peer_single_mode(peer, sys_cfg):
return False
if peer_id is None:
return False
if peer_receives_group_tgid(peer, slot, tgid):
return False
pk = bytes_4(int_id(peer_id))
store = sys_cfg.get("_PEER_UA_MULTI_TGS")
if not isinstance(store, dict):
return False
per_peer = store.get(pk)
if not isinstance(per_peer, dict):
return False
slot_set = per_peer.get(int(slot))
return isinstance(slot_set, set) and int(tgid) in slot_set
def register_peer_ua_session(
peer: dict[str, Any],
peer_id: bytes,
@ -203,8 +258,9 @@ def register_peer_ua_session(
*,
now: float | None = None,
) -> None:
"""Track exclusive SINGLE TG for this hotspot (RX / local TX), per-peer OPTIONS TIMER."""
"""Track UA TG for this hotspot (SINGLE=1 exclusive; SINGLE=0 multi-dynamic set)."""
if not peer_single_mode(peer, sys_cfg):
register_peer_ua_multi_tg(peer, peer_id, slot, tgid, sys_cfg)
return
from adn_server.application.report.payloads import resolve_peer_single_and_timer
@ -300,7 +356,7 @@ def clear_peer_ua_sessions(
*,
slot: int | None = None,
) -> None:
"""Clear per-peer SINGLE session (TG 4000, reconnect)."""
"""Clear per-peer UA state (SINGLE session and/or SINGLE=0 multi-dynamic set)."""
pk = bytes_4(int_id(peer_id))
store = sys_cfg.get("_PEER_UA_SESSIONS")
if isinstance(store, dict) and pk in store:
@ -310,6 +366,14 @@ def clear_peer_ua_sessions(
per_peer = store.get(pk)
if isinstance(per_peer, dict):
per_peer.pop(slot, None)
multi = sys_cfg.get("_PEER_UA_MULTI_TGS")
if isinstance(multi, dict) and pk in multi:
if slot is None:
multi.pop(pk, None)
else:
per_peer = multi.get(pk)
if isinstance(per_peer, dict):
per_peer.pop(int(slot), None)
sessions = peer.get("_UA_SESSION")
if isinstance(sessions, dict):
if slot is None:
@ -429,8 +493,9 @@ def peer_should_receive_group_voice(
1. ``SINGLE=1`` with an active session on another TG → deny all other TGs.
2. TG in this peer's OPTIONS static list → allow (when not blocked by SINGLE).
3. Dynamic UA not in static list but owned by this peer's session → allow.
4. No static TG in OPTIONS and no owned dynamic session → deny (no fan-out).
3. ``SINGLE=1``: dynamic UA owned by this peer's exclusive session → allow.
4. ``SINGLE=0``: dynamic UA this peer keyed (multi set) → allow.
5. Otherwise → deny (no fan-out).
A system-wide ACTIVE bridge leg must **not** fan out to every hotspot; that
was the regression when ``bridges`` were passed into this helper.
@ -442,6 +507,8 @@ def peer_should_receive_group_voice(
return True
if _peer_owns_dynamic_ua(peer, slot, tgid, sys_cfg, peer_id=peer_id, now=now):
return True
if peer_owns_multi_dynamic_ua(peer, slot, tgid, sys_cfg, peer_id=peer_id):
return True
return False

@ -87,6 +87,7 @@ def merge_system_config(old_cfg: dict[str, Any], new_cfg: dict[str, Any]) -> dic
"_options_static_apply_fp",
"_default_options",
"_PEER_UA_SESSIONS",
"_PEER_UA_MULTI_TGS",
):
if key in old_cfg:
merged[key] = old_cfg[key]

@ -7,7 +7,6 @@ import struct
from hashlib import pbkdf2_hmac
from typing import Any
from twisted.internet import reactor
from twisted.internet.defer import inlineCallbacks
from twisted.internet.interfaces import IDelayedCall
from twisted.internet.task import LoopingCall
@ -68,7 +67,7 @@ class ProxySelfServiceBridge:
call.start(interval, now=False)
self._loop_calls.append(call)
self._log.info(
"(SELF_SERVICE) DB options at login, send_opts every 10s, "
"(SELF_SERVICE) DB options on PASS= (immediate), send_opts every 10s, "
"clean_tbl every 1h, lst_seen every 2min"
)
@ -111,8 +110,8 @@ class ProxySelfServiceBridge:
mode = data[97:98].decode("utf-8", errors="replace") if len(data) >= 98 else "4"
callsign = data[8:16].rstrip().decode("utf-8", errors="replace")
self._store.ins_conf(int_id(peer_id), peer_id, callsign, host, mode)
self._cancel_opt_timer(peer_id)
self._opt_timers[peer_id] = reactor.callLater(10, self._login_opt, peer_id)
if peer_id in self._mysql_option_peers:
self._fetch_options_now(peer_id)
def _handle_rpto(
self,
@ -139,6 +138,7 @@ class ProxySelfServiceBridge:
)
self._log.info("(SELF_SERVICE) Password stored for: %s", int_id(peer_id))
self._mysql_option_peers.add(peer_id)
self._fetch_options_now(peer_id)
return True
self._mysql_option_peers.discard(peer_id)
self._store.updt_tbl("opt_rcvd", peer_id)
@ -151,6 +151,20 @@ class ProxySelfServiceBridge:
if timer is not None and timer.active():
timer.cancel()
def _fetch_options_now(self, peer_id: bytes) -> None:
"""Load OPTIONS from MySQL and inject RPTO to the server (no legacy 10s delay)."""
self._cancel_opt_timer(peer_id)
if peer_id not in self._mysql_option_peers:
return
if self._use_cases.resolve_client(peer_id) is None:
return
d = self._login_opt(peer_id)
d.addErrback(
lambda f: self._log.warning(
"(SELF_SERVICE) fetch_options_now error: %s", f.getErrorMessage()
)
)
def _inject_rpto(self, peer_id: bytes, options: str | bytes) -> None:
client = self._use_cases.resolve_client(peer_id)
if client is None:

@ -1081,6 +1081,7 @@ class HBPProtocol(DatagramProtocol):
self._on_options_received(self._system, _this_peer["OPTIONS"])
except Exception:
pass
self._push_config_to_monitor()
else:
self.transport.write(b"".join([MSTNAK, _peer_id]), _sockaddr)
logger.info("(%s) Options from Radio ID that is not logged: %s", self._system, int_id(_peer_id))

@ -318,6 +318,27 @@ def test_report_wire_skips_unchanged_dashboard_state() -> None:
assert wire.state_frames(snapshot, force=False) == ()
def test_report_wire_emits_when_runtime_peer_options_change() -> None:
"""RPTO / MySQL inject must push STATE_SND so monitor chips update without reload."""
peer_specs = [(730039101, 4)]
systems = _inject_runtime_systems(peer_specs)
proxy_cfg = _proxy_config()
slots = _peer_slots(peer_specs)
wire = ReportWire()
expanded = expand_inject_proxy_systems(proxy_cfg, systems, slots)
wire.state_frames(expanded, force=True)
assert wire.state_frames(expanded, force=False) == ()
systems["SYSTEM"]["PEERS"][bytes_4(730039101)]["OPTIONS"] = b"TS2=730444,7305;SINGLE=1;"
expanded2 = expand_inject_proxy_systems(proxy_cfg, systems, slots)
frames = wire.state_frames(expanded2, force=False)
assert len(frames) == 1
doc = json.loads(frames[0][1:].decode())
peer_row = doc["ctable"]["MASTERS"]["SYSTEM-4"]["peers"]["730039101"]
assert peer_row["options"] == "TS2=730444,7305;SINGLE=1;"
assert peer_row["ts2_static"] == ["730444", "7305"]
def test_report_wire_bridge_frames_emits_routing_table() -> None:
bridges = {
"52090": [

@ -11,6 +11,7 @@ from adn_server.application.bridge.helpers import (
peer_should_receive_group_voice,
peer_single_blocks_group_voice,
peer_single_blocks_uplink,
register_peer_ua_multi_tg,
register_peer_ua_session,
seed_peer_ua_session_from_status,
)
@ -268,6 +269,55 @@ def test_peer_without_static_tgs_receives_nothing_until_dynamic() -> None:
)
def test_single_zero_dynamic_heard_when_both_peers_keyed() -> None:
"""SINGLE=0: HS1 and HS2 both keyed 7304 → each hears the other's TX on 7304."""
peer_a = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
peer_b = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
sys_cfg = _sys_cfg()
id_a = bytes_4(730039101)
id_b = bytes_4(730039210)
now = 1_000_000.0
register_peer_ua_session(peer_a, id_a, 2, 7304, sys_cfg, now=now)
register_peer_ua_session(peer_b, id_b, 2, 7304, sys_cfg, now=now + 5)
assert peer_should_receive_group_voice(
peer_a, 2, 7304, peer_id=id_a, connected_count=2, sys_cfg=sys_cfg, now=now + 10
)
assert peer_should_receive_group_voice(
peer_b, 2, 7304, peer_id=id_b, connected_count=2, sys_cfg=sys_cfg, now=now + 10
)
def test_single_zero_dynamic_not_heard_without_local_key() -> None:
"""SINGLE=0: dynamic 7304 only reaches peers that keyed it (not the other HS)."""
peer_a = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
peer_b = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
sys_cfg = _sys_cfg()
id_a = bytes_4(730039101)
id_b = bytes_4(730039210)
now = 1_000_000.0
register_peer_ua_session(peer_a, id_a, 2, 7304, sys_cfg, now=now)
assert peer_should_receive_group_voice(
peer_a, 2, 7304, peer_id=id_a, connected_count=2, sys_cfg=sys_cfg, now=now + 10
)
assert not peer_should_receive_group_voice(
peer_b, 2, 7304, peer_id=id_b, connected_count=2, sys_cfg=sys_cfg, now=now + 10
)
def test_single_zero_tg4000_clears_multi_dynamic() -> None:
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=0;"}
sys_cfg = _sys_cfg()
peer_id = _peer_id()
now = 1_000_000.0
register_peer_ua_multi_tg(peer, peer_id, 2, 7304, sys_cfg)
clear_peer_ua_sessions(peer, sys_cfg, peer_id, slot=2)
assert not peer_should_receive_group_voice(
peer, 2, 7304, peer_id=peer_id, connected_count=2, sys_cfg=sys_cfg, now=now + 10
)
def test_new_tx_replaces_single_session_tg() -> None:
peer = {"OPTIONS": b"TS2=730,7305;SINGLE=1;TIMER=5;"}
sys_cfg = _sys_cfg()

@ -124,7 +124,8 @@ def test_rptc_login_opt_skips_without_pass() -> None:
assert sink.injected == []
def test_rptc_login_opt_injects_rpto_after_pass() -> None:
def test_rptc_after_pass_reinjects_rpto() -> None:
"""RPTC after PASS= re-fetches OPTIONS once ins_conf has run (logged_in in DB)."""
bridge, sink, store, _sender = _bridge()
peer = bytes_4(7300444)
store.options_by_peer[peer] = "TS2=730444;"
@ -133,11 +134,12 @@ def test_rptc_login_opt_injects_rpto_after_pass() -> None:
("192.168.1.10", 62031),
peer,
)
assert len(sink.injected) == 1
packet = RPTC + peer + b"CE1ILI " + b"\x00" * 85 + b"4"
bridge.before_inject(packet, ("192.168.1.10", 62031), peer)
_run_deferred(bridge._login_opt(peer))
assert sink.injected
assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;"
assert ("ins_conf", peer) in store.actions
assert len(sink.injected) == 2
assert sink.injected[-1][0] == RPTO + peer + b"TS2=730444;"
def test_rpto_pass_stores_password_and_skips_inject() -> None:
@ -151,6 +153,23 @@ def test_rpto_pass_stores_password_and_skips_inject() -> None:
assert sender.sent and sender.sent[0][0][:6] == b"RPTACK"
def test_rpto_pass_fetches_options_immediately() -> None:
"""PASS= triggers MySQL slct_opt + RPTO inject without waiting for RPTC/10s timer."""
bridge, sink, store, sender = _bridge()
peer = bytes_4(7300444)
store.options_by_peer[peer] = "TS2=730444;SINGLE=1;"
skip = bridge.before_inject(
RPTO + peer + b"PASS=secret123",
("192.168.1.10", 62031),
peer,
)
assert skip is True
assert ("psswd", peer) in store.actions
assert sender.sent and sender.sent[0][0][:6] == b"RPTACK"
assert sink.injected
assert sink.injected[0][0] == RPTO + peer + b"TS2=730444;SINGLE=1;"
def test_send_opts_skips_without_pass() -> None:
bridge, sink, store, _sender = _bridge()
peer = bytes_4(7300444)

Loading…
Cancel
Save

Powered by TurnKey Linux.