diff --git a/src/adn_server/infrastructure/subscription_store.py b/src/adn_server/infrastructure/subscription_store.py index 19f876c..116696c 100644 --- a/src/adn_server/infrastructure/subscription_store.py +++ b/src/adn_server/infrastructure/subscription_store.py @@ -45,6 +45,9 @@ class InMemorySubscriptionStore(SubscriptionStore): self._by_table: dict[str, list[Subscription]] = defaultdict(list) self._source_tables: dict[_IndexKey, set[str]] = {} self._active_target_counts: dict[_IndexKey, int] = {} + # What each leg was indexed under. Callers change a leg in place and then + # upsert it, so by the time it is unindexed it may no longer say where it is. + self._indexed: dict[SubscriptionId, tuple[str, _IndexKey | None]] = {} def get(self, sub_id: SubscriptionId) -> Subscription | None: return self._items.get(sub_id) @@ -68,6 +71,7 @@ class InMemorySubscriptionStore(SubscriptionStore): self._by_table.clear() self._source_tables.clear() self._active_target_counts.clear() + self._indexed.clear() def replace_all(self, subscriptions: Sequence[Subscription]) -> None: self.clear() @@ -116,13 +120,18 @@ class InMemorySubscriptionStore(SubscriptionStore): def _index(self, sub: Subscription) -> None: table_key = sub.table_key() self._by_table[table_key].append(sub) + key = None if sub.is_active(): key = self._index_key(sub) self._source_tables.setdefault(key, set()).add(table_key) self._active_target_counts[key] = self._active_target_counts.get(key, 0) + 1 + self._indexed[sub.subscription_id] = (table_key, key) def _unindex(self, sub: Subscription) -> None: - table_key = sub.table_key() + indexed = self._indexed.pop(sub.subscription_id, None) + if indexed is None: + return + table_key, key = indexed legs = self._by_table.get(table_key) if legs: try: @@ -131,8 +140,7 @@ class InMemorySubscriptionStore(SubscriptionStore): pass if not legs: del self._by_table[table_key] - if sub.is_active(): - key = self._index_key(sub) + if key is not None: keys = self._source_tables.get(key) if keys is not None: keys.discard(table_key) diff --git a/tests/infrastructure/test_subscription_store_index.py b/tests/infrastructure/test_subscription_store_index.py index 01c8a03..7b60bc8 100644 --- a/tests/infrastructure/test_subscription_store_index.py +++ b/tests/infrastructure/test_subscription_store_index.py @@ -86,3 +86,50 @@ def test_resolve_uses_table_index() -> None: ) assert len(legs) == 1 assert legs[0].target_system == "SYSTEM" + + +# Timers, in-band signalling and resets change a leg in place and then upsert +# the same object, so the store can no longer read the old state off it. + + +def test_in_place_deactivation_leaves_the_active_indexes() -> None: + store = InMemorySubscriptionStore() + store.replace_all(subscriptions_from_routing_table({"7147": [_row(system="SYSTEM", ts=2, tgid=7147, active=True)]})) + sub = store.snapshot()[0] + sub.state.phase = SubscriptionPhase.IDLE + store.upsert(sub) + assert store.relay_tables_with_active_source("SYSTEM", 2, 7147) == () + assert store.has_active_target_leg("SYSTEM", 2, 7147) is False + assert store.legs_in_table("7147") == (sub,) + + +def test_in_place_activation_keeps_other_active_legs_counted() -> None: + # Two legs of SYSTEM delivering on TS2 TG 7147: its own table and a + # translated one (table 9000 rewritten to 7147) share the active index key. + bridges = { + "7147": [_row(system="SYSTEM", ts=2, tgid=7147, active=True)], + "9000": [_row(system="SYSTEM", ts=2, tgid=7147, active=False)], + } + store = InMemorySubscriptionStore() + store.replace_all(subscriptions_from_routing_table(bridges)) + translated = next(s for s in store.snapshot() if s.table_key() == "9000") + translated.state.phase = SubscriptionPhase.ACTIVE + store.upsert(translated) + assert store.relay_tables_with_active_source("SYSTEM", 2, 7147) == ("7147", "9000") + + main = next(s for s in store.snapshot() if s.table_key() == "7147") + main.state.phase = SubscriptionPhase.IDLE + store.upsert(main) + assert store.relay_tables_with_active_source("SYSTEM", 2, 7147) == ("9000",) + assert store.has_active_target_leg("SYSTEM", 2, 7147) is True + + +def test_remove_after_an_in_place_change_clears_every_index() -> None: + store = InMemorySubscriptionStore() + store.replace_all(subscriptions_from_routing_table({"7147": [_row(system="SYSTEM", ts=2, tgid=7147, active=True)]})) + sub = store.snapshot()[0] + sub.state.phase = SubscriptionPhase.IDLE + assert store.remove(sub.subscription_id) is True + assert store.relay_tables_with_active_source("SYSTEM", 2, 7147) == () + assert store.has_active_target_leg("SYSTEM", 2, 7147) is False + assert store.legs_in_table("7147") == () diff --git a/tests/routing/test_rule_timer_stops_source.py b/tests/routing/test_rule_timer_stops_source.py new file mode 100644 index 0000000..a0a8dc7 --- /dev/null +++ b/tests/routing/test_rule_timer_stops_source.py @@ -0,0 +1,63 @@ +# ADN DMR Peer Server - tests routing rule timer stops a timed-out source +# +# Copyright (C) 2026 Rodrigo Pérez, CE5RPY +# +############################################################################### +# This program is free software; you can redistribute it and/or modify +# it under the terms of the GNU General Public License as published by +# the Free Software Foundation; either version 3 of the License, or +# (at your option) any later version. +# +# This program is distributed in the hope that it will be useful, +# but WITHOUT ANY WARRANTY; without even the implied warranty of +# MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the +# GNU General Public License for more details. +# +# You should have received a copy of the GNU General Public License +# along with this program; if not, write to the Free Software Foundation, +# Inc., 51 Franklin Street, Fifth Floor, Boston, MA 02110-1301 USA +############################################################################### + +"""A source leg the rule timer deactivates stops being forwarded (legacy: ACTIVE per frame).""" + +from __future__ import annotations + +import time + +import pytest + +from tests.harness.deterministic import ( + DeterministicScenario, + PacketSpec, + active_routing_table, + minimal_config, +) + + +@pytest.mark.behavior +def test_frames_after_the_source_times_out_are_not_forwarded(monkeypatch: pytest.MonkeyPatch) -> None: + config = minimal_config(("MASTER-A", "MASTER-B")) + config["SYSTEMS"]["MASTER-B"]["SINGLE_MODE"] = False # its leg has no timer + sc = DeterministicScenario( + config=config, + routing_table=active_routing_table(213, (("MASTER-A", 2), ("MASTER-B", 2)), timeout_minutes=1), + ) + monkeypatch.setattr(time, "time", sc.clock.time) + sc.routing.apply_startup_subscriptions() + + base = PacketSpec(dst_id=213, stream_id=0x1234, slot=2) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_head_spec(base)) + for seq in range(1, 4): + sc.clock.advance(0.06) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=seq)) + assert len(sc.capture.for_system("MASTER-B")) == 4 + + sc.clock.advance(120) + sc.routing.rule_timer_loop() + assert sc.subscription_store.relay_tables_with_active_source("MASTER-A", 2, 213) == () + + sc.capture.packets.clear() + for seq in range(4, 8): + sc.clock.advance(0.06) + sc.inject_hbp("MASTER-A", DeterministicScenario.voice_burst_spec(base, seq=seq, dtype_vseq=seq % 6)) + assert sc.capture.for_system("MASTER-B") == []