release: adn-server 2.0.0-alpha1 (pre-GA baseline)

Populate SubscriptionStore from BRIDGES on OPTIONS/startup paths (P2-008).
BRIDGES remains voice routing authority; no dmrd_received changes.

Phases 0-4 complete. Next: P2-009 subscription runtime cutover.
pull/1/head
Rodrigo Pérez 4 months ago
parent cd836d1f5a
commit 67710558e8

@ -7,7 +7,7 @@ build-backend = "setuptools.build_meta"
[project]
name = "adn-server"
version = "1.0.0"
version = "2.0.0-alpha1"
description = "ADN DMR Peer Server"
readme = "README.md"
license = { text = "GPL-3.0-or-later" }

@ -37,4 +37,4 @@
"""ADN DMR Peer Server — conference bridge (rewrite of bridge_master)."""
__version__ = "1.0.0"
__version__ = "2.0.0-alpha1"

@ -314,6 +314,7 @@ class BridgeTableMixin:
if tg in prohibited_tgs:
continue
self.make_static_tg(tg, 2, tmout, system)
self._sync_subscription_store()
def options_config_for_system(self, system_name: str) -> None:
"""Update static TG bridges for one system immediately (e.g. when RPTO received). So incoming OBP traffic reaches hotspots without waiting for the 26s options_config_loop."""
@ -379,6 +380,7 @@ class BridgeTableMixin:
_fp = f"{new_ts1}|{new_ts2}|{int(_tmout)}"
if sys_cfg.get("_options_static_apply_fp") == _fp:
self._restore_prohibited_static_bridge_legs(system_name)
self._sync_subscription_store()
return
# Legacy: reset TGs that were removed (bridge_master.py 1736-1767)
old_ts1 = str(sys_cfg.get("TS1_STATIC") or "").strip()
@ -444,6 +446,7 @@ class BridgeTableMixin:
if new_ts1 or new_ts2:
logger.info("(OPTIONS) %s static TGs applied: TS1=%s TS2=%s", system_name, new_ts1 or "-", new_ts2 or "-")
self._restore_prohibited_static_bridge_legs(system_name)
self._sync_subscription_store()
except Exception as e:
logger.debug("(OPTIONS) options_config_for_system %s: %s", system_name, e)
@ -827,4 +830,5 @@ class BridgeTableMixin:
systems_cfg[_system]["DEFAULT_UA_TIMER"] = int(_options.get("DEFAULT_UA_TIMER", 10))
except Exception as e:
logger.exception("(OPTIONS) caught exception: %s", e)
self._sync_subscription_store()

@ -40,7 +40,8 @@ from bitarray import bitarray
from ..domain.dmr import bptc
from ..domain import HBPF_DATA_SYNC, HBPF_SLT_VHEAD, HBPF_SLT_VTERM, STREAM_TO, bytes_3, bytes_4, int_id
from .ports import BridgeRouter, DmrEmbeddedLcEncoder, TalkerAliasEmblcEncoder
from .ports import BridgeRouter, DmrEmbeddedLcEncoder, SubscriptionStore, TalkerAliasEmblcEncoder
from .subscription.store_sync import replace_store_from_bridges
from .talker_alias_use_cases import TalkerAliasUseCases
from .bridge.helpers import obp_target_bcsq_quenches_stream, resolve_voice_peer_id
from .bridge.timers import BridgeTimerMixin
@ -69,9 +70,11 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi
call_later: Any = None,
encode_emblc: DmrEmbeddedLcEncoder | None = None,
ta_emblc_encoder: TalkerAliasEmblcEncoder | None = None,
subscription_store: SubscriptionStore | None = None,
) -> None:
self._router = bridge_router
self._config = config
self._subscription_store = subscription_store
self._send_to_system = send_to_system # (system_name, packet, **kwargs) -> None
self._get_protocols = get_protocols # () -> dict[str, protocol]
self._reporting = reporting
@ -101,6 +104,12 @@ class BridgeUseCases(BridgeTimerMixin, BridgeObpForwardMixin, BridgeHbpForwardMi
except Exception:
return False
def _sync_subscription_store(self) -> None:
"""Mirror BRIDGES into the subscription store (OPTIONS/static paths only for now)."""
if self._subscription_store is None:
return
replace_store_from_bridges(self._subscription_store, self._router.get_bridges())
def get_bridges(self) -> dict[str, list[dict[str, Any]]]:
"""Return current BRIDGES."""
return self._router.get_bridges()

@ -1,6 +1,8 @@
"""Subscription application helpers."""
from .bridges_export import export_bridges, subscription_to_legacy_row
from .bridges_import import subscriptions_from_bridges
from .store_sync import replace_store_from_bridges
from .bridges_legacy_view import BridgesLegacyView
from .ingress import build_voice_ingress
from .router import SubscriptionRouter
@ -10,5 +12,7 @@ __all__ = [
"SubscriptionRouter",
"build_voice_ingress",
"export_bridges",
"replace_store_from_bridges",
"subscription_to_legacy_row",
"subscriptions_from_bridges",
]

@ -0,0 +1,75 @@
"""Import legacy ``BRIDGES`` rows into domain subscriptions (mirror of ``bridges_export``)."""
from __future__ import annotations
from typing import Any
from adn_server.domain import int_id
from adn_server.domain.subscription import (
ActivationPolicy,
AudioChannel,
Subscription,
SubscriptionPhase,
SubscriptionRole,
SubscriptionState,
SystemId,
TgId,
)
def subscriptions_from_bridges(bridges: dict[str, list[dict[str, Any]]]) -> list[Subscription]:
"""Build subscriptions from a legacy ``BRIDGES`` snapshot (OPTIONS/static TG / ECHO)."""
subs: list[Subscription] = []
for table_key, rows in bridges.items():
for row in rows:
if not isinstance(row, dict):
continue
subs.append(_subscription_from_row(table_key, row))
return subs
def _subscription_from_row(table_key: str, row: dict[str, Any]) -> Subscription:
ts = int(row.get("TS") or 1)
if table_key.startswith("#"):
channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00")
bridge_key = table_key
else:
try:
channel_tgid = int(table_key)
except ValueError:
channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00")
bridge_key = None
to_type = str(row.get("TO_TYPE", "ON"))
timer = row.get("TIMER")
timer_at = float(timer) if isinstance(timer, (int, float)) else None
timeout = row.get("TIMEOUT")
timeout_sec = float(timeout) if isinstance(timeout, (int, float)) else None
return Subscription(
channel=AudioChannel(tgid=TgId(channel_tgid), slot=ts), # type: ignore[arg-type]
system=SystemId(str(row.get("SYSTEM", ""))),
target_tgid=TgId(int_id(row.get("TGID") or b"\x00\x00\x00")),
role=_role_from_to_type(to_type),
policy=_policy_from_to_type(to_type),
state=SubscriptionState(
phase=SubscriptionPhase.ACTIVE if row.get("ACTIVE") else SubscriptionPhase.IDLE,
timer_expires_at=timer_at,
),
bridge_key=bridge_key,
timeout_seconds=timeout_sec,
)
def _role_from_to_type(to_type: str) -> SubscriptionRole:
if to_type == "NONE":
return SubscriptionRole.ECHO
if to_type == "STAT":
return SubscriptionRole.PASSIVE_STAT
return SubscriptionRole.SINK
def _policy_from_to_type(to_type: str) -> ActivationPolicy:
if to_type == "STAT":
return ActivationPolicy.OPENBRIDGE_STAT
if to_type == "OFF":
return ActivationPolicy.STATIC
return ActivationPolicy.INBAND

@ -0,0 +1,16 @@
"""Keep ``SubscriptionStore`` aligned with legacy ``BRIDGES`` (read-only mirror; not routing authority yet)."""
from __future__ import annotations
from typing import Any
from adn_server.application.ports import SubscriptionStore
from adn_server.application.subscription.bridges_import import subscriptions_from_bridges
def replace_store_from_bridges(
store: SubscriptionStore,
bridges: dict[str, list[dict[str, Any]]],
) -> None:
"""Replace store contents from a full ``BRIDGES`` snapshot."""
store.replace_all(subscriptions_from_bridges(bridges))

@ -36,6 +36,7 @@ from adn_server.application.runtime_context import (
from adn_server.domain import bytes_3
from adn_server.domain.dmr.bptc import encode_emblc
from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
from adn_server.infrastructure.config_normalizer import (
ensure_system_runtime_config as _ensure_system_runtime_config,
)
@ -303,6 +304,7 @@ def run_peer_server(
call_from_reactor=reactor.callFromThread,
)
subscription_store = InMemorySubscriptionStore()
bridge_use_cases = BridgeUseCases(
bridge_router,
config,
@ -316,6 +318,7 @@ def run_peer_server(
call_later=reactor.callLater,
encode_emblc=encode_emblc,
ta_emblc_encoder=default_ta_emblc_encoder,
subscription_store=subscription_store,
)
bridge_use_cases.apply_startup_bridges()
report_factory.set_bridges(bridge_router.get_bridges())

@ -0,0 +1,90 @@
"""BRIDGES → Subscription import (policy / role parity)."""
from __future__ import annotations
from adn_server.application.subscription.bridges_import import subscriptions_from_bridges
from adn_server.domain import bytes_3
from adn_server.domain.subscription import (
ActivationPolicy,
SubscriptionPhase,
SubscriptionRole,
)
def test_import_echo_and_inband_rows():
tgid_b = bytes_3(9990)
bridges = {
"9990": [
{
"SYSTEM": "ECHO",
"TS": 2,
"TGID": tgid_b,
"ACTIVE": True,
"TIMEOUT": 120.0,
"TO_TYPE": "NONE",
},
{
"SYSTEM": "MASTER-A",
"TS": 1,
"TGID": tgid_b,
"ACTIVE": False,
"TIMEOUT": 600.0,
"TO_TYPE": "ON",
},
]
}
subs = subscriptions_from_bridges(bridges)
assert len(subs) == 2
echo, master = subs
assert echo.role == SubscriptionRole.ECHO
assert echo.policy == ActivationPolicy.INBAND
assert echo.state.phase == SubscriptionPhase.ACTIVE
assert master.role == SubscriptionRole.SINK
assert master.policy == ActivationPolicy.INBAND
assert master.state.phase == SubscriptionPhase.IDLE
def test_import_static_and_stat_rows():
bridges = {
"12345": [
{
"SYSTEM": "MASTER-A",
"TS": 1,
"TGID": bytes_3(12345),
"ACTIVE": True,
"TIMEOUT": 600.0,
"TO_TYPE": "OFF",
},
{
"SYSTEM": "OBP-CL",
"TS": 1,
"TGID": bytes_3(12345),
"ACTIVE": True,
"TIMEOUT": "",
"TO_TYPE": "STAT",
},
]
}
subs = subscriptions_from_bridges(bridges)
static, stat = subs
assert static.policy == ActivationPolicy.STATIC
assert static.role == SubscriptionRole.SINK
assert stat.policy == ActivationPolicy.OPENBRIDGE_STAT
assert stat.role == SubscriptionRole.PASSIVE_STAT
def test_import_hash_table_key_sets_bridge_key():
bridges = {
"#730444": [
{
"SYSTEM": "OBP-CL",
"TS": 1,
"TGID": bytes_3(730444),
"ACTIVE": True,
"TIMEOUT": "",
"TO_TYPE": "STAT",
}
]
}
(sub,) = subscriptions_from_bridges(bridges)
assert sub.bridge_key == "#730444"

@ -0,0 +1,84 @@
"""SubscriptionStore sync from BRIDGES (V2-P2-008)."""
from __future__ import annotations
from typing import Any
from adn_server.application.subscription.bridges_export import export_bridges
from adn_server.application.subscription.store_sync import replace_store_from_bridges
from adn_server.domain import bytes_3, int_id
from adn_server.infrastructure.bootstrap.peer_server import _make_echo_bridges
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _row_fingerprint(row: dict[str, Any]) -> tuple[Any, ...]:
return (
row.get("SYSTEM"),
int(row.get("TS") or 1),
int_id(row.get("TGID") or b"\x00\x00\x00"),
bool(row.get("ACTIVE")),
str(row.get("TO_TYPE", "ON")),
row.get("TIMEOUT"),
)
def _bridges_fingerprints(bridges: dict[str, list[dict[str, Any]]]) -> set[tuple[Any, ...]]:
fps: set[tuple[Any, ...]] = set()
for rows in bridges.values():
for row in rows:
if isinstance(row, dict):
fps.add(_row_fingerprint(row))
return fps
def test_replace_store_from_echo_bridges_round_trip():
config = {
"SYSTEMS": {
"ECHO": {"MODE": "PEER"},
"MASTER-A": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 10},
"MASTER-B": {"MODE": "MASTER", "DEFAULT_UA_TIMER": 15},
}
}
bridges = _make_echo_bridges(config)
store = InMemorySubscriptionStore()
replace_store_from_bridges(store, bridges)
assert len(store.snapshot()) == len(_bridges_fingerprints(bridges))
now = 1_700_000_000.0
exported = export_bridges(store, now=now)
assert _bridges_fingerprints(bridges) == _bridges_fingerprints(exported)
def test_replace_store_clears_stale_entries():
store = InMemorySubscriptionStore()
bridges_a = {
"100": [
{
"SYSTEM": "SYS-A",
"TS": 1,
"TGID": bytes_3(100),
"ACTIVE": False,
"TIMEOUT": 600.0,
"TO_TYPE": "ON",
}
]
}
bridges_b = {
"200": [
{
"SYSTEM": "SYS-B",
"TS": 2,
"TGID": bytes_3(200),
"ACTIVE": True,
"TIMEOUT": 600.0,
"TO_TYPE": "ON",
}
]
}
replace_store_from_bridges(store, bridges_a)
assert len(store.snapshot()) == 1
replace_store_from_bridges(store, bridges_b)
assert len(store.snapshot()) == 1
(sub,) = store.snapshot()
assert int(sub.channel.tgid) == 200

@ -4,61 +4,15 @@ from __future__ import annotations
from typing import Any
from adn_server.application.subscription.bridges_import import subscriptions_from_bridges
from adn_server.application.subscription.router import SubscriptionRouter
from adn_server.domain import bytes_3, int_id
from adn_server.domain.subscription import (
ActivationPolicy,
AudioChannel,
Subscription,
SubscriptionPhase,
SubscriptionRole,
SubscriptionState,
SystemId,
TgId,
)
from adn_server.domain.subscription import TgId
from adn_server.domain.voice_routing import ForwardLeg, VoiceIngress
from adn_server.infrastructure.bridge_router_impl import InMemoryBridgeRouter
from adn_server.infrastructure.subscription_store import InMemorySubscriptionStore
def _role_from_to_type(to_type: str) -> SubscriptionRole:
if to_type == "NONE":
return SubscriptionRole.ECHO
if to_type == "STAT":
return SubscriptionRole.PASSIVE_STAT
return SubscriptionRole.SINK
def subscriptions_from_bridges(bridges: dict[str, list[dict[str, Any]]]) -> list[Subscription]:
"""Build subscriptions that round-trip ``export_bridges`` shape for parity tests."""
subs: list[Subscription] = []
for table_key, rows in bridges.items():
for row in rows:
if not isinstance(row, dict):
continue
ts = int(row.get("TS") or 1)
if table_key.startswith("#"):
channel_tgid = int_id(row.get("TGID") or b"\x00\x00\x00")
else:
channel_tgid = int(table_key)
to_type = str(row.get("TO_TYPE", "ON"))
subs.append(
Subscription(
channel=AudioChannel(tgid=TgId(channel_tgid), slot=ts), # type: ignore[arg-type]
system=SystemId(str(row.get("SYSTEM", ""))),
target_tgid=TgId(int_id(row.get("TGID") or b"\x00\x00\x00")),
role=_role_from_to_type(to_type),
policy=ActivationPolicy.INBAND,
state=SubscriptionState(
phase=SubscriptionPhase.ACTIVE if row.get("ACTIVE") else SubscriptionPhase.IDLE
),
bridge_key=table_key if table_key.startswith("#") else None,
timeout_seconds=row.get("TIMEOUT") if isinstance(row.get("TIMEOUT"), (int, float)) else None,
)
)
return subs
def legacy_forward_targets(
bridges: dict[str, list[dict[str, Any]]],
*,

Loading…
Cancel
Save

Powered by TurnKey Linux.