fix: hard one-QSO slot rule for all peers and stale session cleanup

Remove the lab-witness exemption so every hotspot obeys one group QSO per RF
slot. Expire same-TG zombie peer_voice_slots after STREAM_TO, allow orphan
VTERM through the slot gate, and route ingress slot tracking via
track_peer_group_dmrd for consistent session lifecycle.
pull/29/head
Rodrigo Pérez 3 months ago
parent de385660c1
commit e87332a1d6

@ -41,6 +41,7 @@ from .helpers import (
parse_dmrd_route_fields, parse_dmrd_route_fields,
peer_downlink_voice_slot, peer_downlink_voice_slot,
peer_is_simplex, peer_is_simplex,
peer_options_static_tg_slot,
peer_receives_group_tgid, peer_receives_group_tgid,
peer_should_receive_group_voice, peer_should_receive_group_voice,
peer_single_exclusive_tgid, peer_single_exclusive_tgid,
@ -115,12 +116,6 @@ def peer_listen_slots(peer: dict[str, Any], tgid: int) -> list[int]:
return [] return []
def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None:
from adn_server.application.routing.helpers import peer_options_static_tg_slot as _slot
return _slot(peer, tgid)
def peer_accepts_group_downlink( def peer_accepts_group_downlink(
ctx: DownlinkContext, ctx: DownlinkContext,
peer_id: bytes, peer_id: bytes,
@ -217,9 +212,11 @@ def peer_slot_blocks_downlink(
if not rx_listening: if not rx_listening:
return True return True
active = per_slot.get(int(listen_slot)) active = per_slot.get(int(listen_slot))
if not isinstance(active, dict) or active.get("stream_id") != stream_id: if not isinstance(active, dict):
return True return False
return False if active.get("stream_id") == stream_id:
return False
return True
if not ctx.per_peer_contention(): if not ctx.per_peer_contention():
return False return False
hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0) hang = float(ctx.sys_cfg.get("GROUP_HANGTIME", 0) or 0)
@ -357,15 +354,19 @@ def track_peer_group_dmrd(
*, *,
pkt_time: float | None = None, pkt_time: float | None = None,
from_ingress: bool = False, from_ingress: bool = False,
voice_slot: int | None = None,
) -> None: ) -> None:
"""Update per-hotspot slot state from downlink or ingress DMRD.""" """Update per-hotspot slot state from downlink or ingress DMRD."""
burst = parse_dmrd_burst_fields(packet) burst = parse_dmrd_burst_fields(packet)
if burst is None: if burst is None:
return return
wire_slot, frame_type, dtype_vseq, stream_id, dst_id, _call_type = burst wire_slot, frame_type, dtype_vseq, stream_id, dst_id, _call_type = burst
voice_slot = peer_downlink_voice_slot( if voice_slot is None:
peer, wire_slot, int_id(dst_id), ctx.sys_cfg, peer_id=peer_id, voice_slot = peer_downlink_voice_slot(
) peer, wire_slot, int_id(dst_id), ctx.sys_cfg, peer_id=peer_id,
)
else:
voice_slot = int(voice_slot)
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM:
if ( if (
not from_ingress not from_ingress
@ -436,9 +437,23 @@ def track_peer_group_dmrd(
) )
def build_dmra_route_packet(slot: int, tgid: int, stream_id: bytes | None = None) -> bytes: def peer_accepts_group_dmrd_packet(
"""Synthetic DMRD for DMRA / monitor fan-out lookup.""" ctx: DownlinkContext,
return synthetic_group_dmrd_route_packet(slot, tgid, stream_id) peer_id: bytes,
peer: dict[str, Any],
packet: bytes,
*,
routed: bool = False,
) -> bool:
"""True when group/vcsbk DMRD passes per-hotspot slot gate (OPTIONS checked separately)."""
if not ctx.per_peer_contention():
return True
route_pkt = (
packet
if routed
else remap_dmrd_for_peer(packet, peer, ctx.sys_cfg, peer_id=peer_id)
)
return not peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt)
def peer_accepts_dmra( def peer_accepts_dmra(
@ -451,7 +466,7 @@ def peer_accepts_dmra(
peer = ctx.peers.get(peer_id) peer = ctx.peers.get(peer_id)
if not isinstance(peer, dict): if not isinstance(peer, dict):
return False return False
route_pkt = build_dmra_route_packet(slot, tgid) route_pkt = synthetic_group_dmrd_route_packet(slot, tgid)
if not peer_accepts_group_downlink(ctx, peer_id, peer, slot, tgid): if not peer_accepts_group_downlink(ctx, peer_id, peer, slot, tgid):
return False return False
remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id) remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id)
@ -477,7 +492,7 @@ def peer_would_show_group_voice_on_monitor(
return False return False
if ctx is None or not ctx.per_peer_contention(): if ctx is None or not ctx.per_peer_contention():
return True return True
route_pkt = build_dmra_route_packet(wire_slot, tgid, stream_id) route_pkt = synthetic_group_dmrd_route_packet(wire_slot, tgid, stream_id)
remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id) remapped = remap_dmrd_for_peer(route_pkt, peer, ctx.sys_cfg, peer_id=peer_id)
return not peer_slot_blocks_downlink(ctx, peer_id, peer, remapped) return not peer_slot_blocks_downlink(ctx, peer_id, peer, remapped)

@ -271,6 +271,31 @@ def _peer_status_rx_hangtime_blocks(
return (pkt_time - rx_t) < hang return (pkt_time - rx_t) < hang
def _peer_single_locked_tgid(
peer: dict[str, Any],
sys_cfg: dict[str, Any],
*,
peer_id: bytes | None = None,
now: float | None = None,
prefer_slot: int | None = None,
) -> int | None:
"""First active SINGLE=1 session TG (``prefer_slot`` checked before the other TS)."""
slots: list[int] = []
if prefer_slot is not None:
slots.append(int(prefer_slot))
for voice_slot in (1, 2):
if prefer_slot is not None and voice_slot == int(prefer_slot):
continue
slots.append(voice_slot)
for voice_slot in slots:
locked = peer_single_exclusive_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now,
)
if locked is not None:
return locked
return None
def peer_single_same_tg_foreign_tx_blocks( def peer_single_same_tg_foreign_tx_blocks(
peer: dict[str, Any], peer: dict[str, Any],
peer_id: bytes, peer_id: bytes,
@ -285,13 +310,9 @@ def peer_single_same_tg_foreign_tx_blocks(
if not sys_cfg or not peer_single_mode(peer, sys_cfg): if not sys_cfg or not peer_single_mode(peer, sys_cfg):
return False return False
incoming = int_id(incoming_tgid_b) incoming = int_id(incoming_tgid_b)
locked = None locked = _peer_single_locked_tgid(
for voice_slot in (1, 2): peer, sys_cfg, peer_id=peer_id, now=pkt_time,
locked = peer_single_exclusive_tgid( )
peer, voice_slot, sys_cfg, peer_id=peer_id, now=pkt_time,
)
if locked is not None:
break
if locked is None or int(locked) != incoming: if locked is None or int(locked) != incoming:
return False return False
if not slot_has_active_voice(slot_st, pkt_time): if not slot_has_active_voice(slot_st, pkt_time):
@ -339,10 +360,22 @@ def peer_hotspot_voice_slot_busy(
if isinstance(active, dict): if isinstance(active, dict):
if active.get("ingress"): if active.get("ingress"):
return True return True
if peer_hotspot_hard_slot_one_qso(peer): active_stream = active.get("stream_id")
active_stream = active.get("stream_id") incoming_tgid = int_id(incoming_tgid_b)
if not (stream_id and active_stream and active_stream == stream_id): active_tgid = int(active.get("tgid", 0) or 0)
return True if stream_id and active_stream and active_stream == stream_id:
pass
elif (
stream_id
and active_stream
and active_stream != stream_id
and active_tgid
and active_tgid == incoming_tgid
and (pkt_time - float(active.get("time", 0) or 0)) >= STREAM_TO
):
peer_slots.pop(int(voice_slot), None)
else:
return True
if isinstance(peer, dict) and peer_single_blocks_foreign_same_tg_downlink( if isinstance(peer, dict) and peer_single_blocks_foreign_same_tg_downlink(
peer, pk, voice_slot, incoming_tgid_b, peer_slots, sys_cfg, now=pkt_time, peer, pk, voice_slot, incoming_tgid_b, peer_slots, sys_cfg, now=pkt_time,
): ):
@ -1302,18 +1335,9 @@ def peer_single_blocks_foreign_same_tg_downlink(
if not sys_cfg or not peer_single_mode(peer, sys_cfg): if not sys_cfg or not peer_single_mode(peer, sys_cfg):
return False return False
incoming = int_id(incoming_tgid_b) incoming = int_id(incoming_tgid_b)
locked = peer_single_exclusive_tgid( locked = _peer_single_locked_tgid(
peer, voice_slot, sys_cfg, peer_id=peer_id, now=now, peer, sys_cfg, peer_id=peer_id, now=now, prefer_slot=voice_slot,
) )
if locked is None:
for alt_slot in (1, 2):
if alt_slot == voice_slot:
continue
locked = peer_single_exclusive_tgid(
peer, alt_slot, sys_cfg, peer_id=peer_id, now=now,
)
if locked is not None:
break
if locked is None or int(locked) != incoming: if locked is None or int(locked) != incoming:
return False return False
active = (peer_slots or {}).get(int(voice_slot)) active = (peer_slots or {}).get(int(voice_slot))
@ -1340,13 +1364,6 @@ def peer_wants_downlink_single_listen_lock(peer: dict[str, Any], sys_cfg: dict[s
return peer_static_options_tg_count(peer) <= 6 return peer_static_options_tg_count(peer) <= 6
def peer_hotspot_hard_slot_one_qso(peer: dict[str, Any] | None) -> bool:
"""True when the one-QSO-per-RF-slot rule applies (normal hotspots, not lab witnesses)."""
if not isinstance(peer, dict):
return True
return peer_static_options_tg_count(peer) <= 6
def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool: def peer_receives_group_tgid(peer: dict[str, Any], slot: int, tgid: int) -> bool:
"""True when peer RPTO OPTIONS list the group TG on TS1 or TS2 (legacy REPEAT parity). """True when peer RPTO OPTIONS list the group TG on TS1 or TS2 (legacy REPEAT parity).
@ -1378,17 +1395,36 @@ def peer_options_static_tg_slot(peer: dict[str, Any], tgid: int) -> int | None:
return None return None
def synthetic_group_dmrd_route_packet( def synthetic_group_dmrd_burst_packet(
slot: int, slot: int,
tgid: int, tgid: int,
stream_id: bytes | None = None, stream_id: bytes,
*,
frame_type: int = 0,
dtype_vseq: int = 0,
call_type: str = "group",
) -> bytes: ) -> bytes:
"""Minimal DMRD for downlink/monitor gate lookup (slot, TG, optional stream).""" """Minimal DMRD with burst header fields for slot-tracking helpers."""
bits = 0x80 if int(slot) == 2 else 0 bits = 0x80 if int(slot) == 2 else 0
if call_type == "vcsbk":
bits |= 0x23
else:
bits |= (int(frame_type) & 0x3) << 4
bits |= int(dtype_vseq) & 0xF
sid = bytes_4(int_id(stream_id)) if stream_id else b"\x00" * 4 sid = bytes_4(int_id(stream_id)) if stream_id else b"\x00" * 4
return b"DMRD" + b"\x00" * 4 + bytes_3(tgid) + b"\x00" * 4 + bytes([bits]) + sid + b"\x00" * 34 return b"DMRD" + b"\x00" * 4 + bytes_3(tgid) + b"\x00" * 4 + bytes([bits]) + sid + b"\x00" * 34
def synthetic_group_dmrd_route_packet(
slot: int,
tgid: int,
stream_id: bytes | None = None,
) -> bytes:
"""Minimal DMRD for downlink/monitor gate lookup (slot, TG, optional stream)."""
sid = stream_id if stream_id else b"\x00" * 4
return synthetic_group_dmrd_burst_packet(slot, tgid, sid)
def peer_downlink_voice_slot( def peer_downlink_voice_slot(
peer: dict[str, Any], peer: dict[str, Any],
wire_slot: int, wire_slot: int,

@ -41,11 +41,10 @@ from twisted.internet.protocol import DatagramProtocol
from ...application.proxy.deployment import is_proxy_inject_only from ...application.proxy.deployment import is_proxy_inject_only
from ...application.routing.downlink import ( from ...application.routing.downlink import (
DownlinkContext, DownlinkContext,
build_dmra_route_packet,
iter_downlink_voice_slots, iter_downlink_voice_slots,
normalize_ua_voice_slot, normalize_ua_voice_slot,
peer_accepts_dmra, peer_accepts_dmra,
peer_slot_blocks_downlink, peer_accepts_group_dmrd_packet,
remap_dmrd_for_peer, remap_dmrd_for_peer,
track_peer_group_dmrd, track_peer_group_dmrd,
) )
@ -65,6 +64,8 @@ from ...application.routing.helpers import (
register_peer_ua_session, register_peer_ua_session,
resolve_voice_peer_id, resolve_voice_peer_id,
seed_peer_ua_session_from_status, seed_peer_ua_session_from_status,
synthetic_group_dmrd_route_packet,
synthetic_group_dmrd_burst_packet,
tg4000_reset_on_vhead, tg4000_reset_on_vhead,
) )
from ...application.routing.peer_downlink_index import ( from ...application.routing.peer_downlink_index import (
@ -619,24 +620,26 @@ class HBPProtocol(DatagramProtocol):
peer = self._peers.get(peer_id) peer = self._peers.get(peer_id)
if peer is None: if peer is None:
return return
from ...application.routing.downlink import normalize_ua_voice_slot
voice_slot = normalize_ua_voice_slot(peer, wire_slot)
ctx = self._downlink_ctx() ctx = self._downlink_ctx()
if not ctx.per_peer_contention(): if not ctx.per_peer_contention():
return return
if frame_type == HBPF_DATA_SYNC and dtype_vseq == HBPF_SLT_VTERM: packet = synthetic_group_dmrd_burst_packet(
from ...application.routing.downlink import end_peer_voice_slot wire_slot,
int_id(dst_id),
end_peer_voice_slot( stream_id,
ctx, peer_id, voice_slot, stream_id, dst_id, pkt_time=pkt_time, frame_type=frame_type,
) dtype_vseq=dtype_vseq,
else: call_type=call_type,
from ...application.routing.downlink import touch_peer_voice_slot )
track_peer_group_dmrd(
touch_peer_voice_slot( ctx,
ctx, peer_id, voice_slot, stream_id, dst_id, pkt_time=pkt_time, ingress=True, peer_id,
) packet,
peer,
pkt_time=pkt_time,
from_ingress=True,
voice_slot=normalize_ua_voice_slot(peer, wire_slot),
)
def _peer_would_accept_group_dmrd( def _peer_would_accept_group_dmrd(
self, self,
@ -651,15 +654,9 @@ class HBPProtocol(DatagramProtocol):
peer = self._peers.get(peer_id) peer = self._peers.get(peer_id)
if peer is None: if peer is None:
return True return True
ctx = self._downlink_ctx() return peer_accepts_group_dmrd_packet(
if not ctx.per_peer_contention(): self._downlink_ctx(), peer_id, peer, packet, routed=routed,
return True
route_pkt = (
packet
if routed
else remap_dmrd_for_peer(packet, peer, self._config, peer_id=peer_id)
) )
return not peer_slot_blocks_downlink(ctx, peer_id, peer, route_pkt)
def _downlink_drop_key(self, peer_id: bytes, route_pkt: bytes) -> tuple[bytes, bytes]: def _downlink_drop_key(self, peer_id: bytes, route_pkt: bytes) -> tuple[bytes, bytes]:
pk = bytes_4(int_id(peer_id)) pk = bytes_4(int_id(peer_id))
@ -844,7 +841,7 @@ class HBPProtocol(DatagramProtocol):
return 0 return 0
ctx = self._downlink_ctx() ctx = self._downlink_ctx()
route_pkt = build_dmra_route_packet(slot, tgid) route_pkt = synthetic_group_dmrd_route_packet(slot, tgid)
sent = 0 sent = 0
for peer_id in self._iter_downlink_peers(route_pkt): for peer_id in self._iter_downlink_peers(route_pkt):
if exclude_peer and peer_id == exclude_peer: if exclude_peer and peer_id == exclude_peer:

@ -249,8 +249,8 @@ def test_per_peer_obp_tx_stamp_blocked_when_other_stream_active() -> None:
) is True ) is True
def test_lab_witness_many_static_tgs_allows_second_tg_on_slot() -> None: def test_lab_witness_many_static_tgs_blocks_second_tg_on_slot() -> None:
"""Full-table lab witness (>6 static TGs) is not subject to one-QSO hard slot lock.""" """Full-table lab witness (>6 static TGs) still obeys one-QSO-per-RF-slot."""
now = 1_000_000.0 now = 1_000_000.0
witness_id = bytes_4(730039257) witness_id = bytes_4(730039257)
tg_list = ",".join(str(730500 + i) for i in range(13)) tg_list = ",".join(str(730500 + i) for i in range(13))
@ -259,7 +259,7 @@ def test_lab_witness_many_static_tgs_allows_second_tg_on_slot() -> None:
2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 1.0}, 2: {"stream_id": _STREAM_A, "tgid": 730500, "time": now - 1.0},
} }
slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM} slot = {"RX_TYPE": HBPF_SLT_VTERM, "TX_TYPE": HBPF_SLT_VTERM}
assert not peer_hotspot_voice_slot_busy( assert peer_hotspot_voice_slot_busy(
witness_id, witness_id,
2, 2,
_STREAM_B, _STREAM_B,
@ -518,8 +518,8 @@ def test_single0_ingress_tx_blocks_other_static_tg_on_slot() -> None:
) )
def test_lab_witness_nine_static_tgs_allows_second_tg_on_slot() -> None: def test_lab_witness_nine_static_tgs_blocks_second_tg_on_slot() -> None:
"""Lab witness with >6 static TGs is not hard-locked to one QSO per RF slot.""" """Lab witness with >6 static TGs is hard-locked to one QSO per RF slot."""
now = 1_000_000.0 now = 1_000_000.0
hs = bytes_4(730039257) hs = bytes_4(730039257)
peer = { peer = {
@ -536,7 +536,7 @@ def test_lab_witness_nine_static_tgs_allows_second_tg_on_slot() -> None:
peer_slots = { peer_slots = {
2: {"stream_id": bytes_4(0x11111111), "tgid": 730507, "time": now, "ingress": False}, 2: {"stream_id": bytes_4(0x11111111), "tgid": 730507, "time": now, "ingress": False},
} }
assert not peer_hotspot_voice_slot_busy( assert peer_hotspot_voice_slot_busy(
hs, hs,
2, 2,
bytes_4(0x22222222), bytes_4(0x22222222),

Loading…
Cancel
Save

Powered by TurnKey Linux.