refactor: extract ocf/ — fold sensors.index_links into StateCache.index_device_tree
This commit is contained in:
@@ -20,7 +20,7 @@ from __future__ import annotations
|
||||
import threading
|
||||
from typing import Callable, Optional
|
||||
|
||||
from .coap_dtls import DtlsCoapSession
|
||||
from protocol.dtls_session import DtlsCoapSession
|
||||
|
||||
|
||||
class KeepaliveTask:
|
||||
@@ -17,7 +17,7 @@ from __future__ import annotations
|
||||
import threading
|
||||
from typing import Optional
|
||||
|
||||
from .coap_dtls import DtlsCoapSession
|
||||
from protocol.dtls_session import DtlsCoapSession
|
||||
|
||||
|
||||
class ObserveRefreshTask:
|
||||
@@ -36,7 +36,7 @@ from typing import Callable, Optional, TYPE_CHECKING
|
||||
|
||||
import cbor2
|
||||
|
||||
from .coap_dtls import DtlsCoapSession, fmt_code
|
||||
from protocol.dtls_session import DtlsCoapSession, fmt_code
|
||||
|
||||
if TYPE_CHECKING:
|
||||
from .state_cache import StateCache
|
||||
@@ -61,7 +61,6 @@ class PollScheduler:
|
||||
session: DtlsCoapSession,
|
||||
cache: 'StateCache',
|
||||
tiers: list[PollTier],
|
||||
sweep_index_fn: Callable[[object], dict[str, dict]],
|
||||
is_active_fn: Optional[Callable[[dict[str, dict]], bool]] = None,
|
||||
logger=None,
|
||||
timeout_s: float = 8.0,
|
||||
@@ -69,7 +68,6 @@ class PollScheduler:
|
||||
self.session = session
|
||||
self.cache = cache
|
||||
self.tiers = tiers
|
||||
self.sweep_index = sweep_index_fn
|
||||
self.is_active_fn = is_active_fn
|
||||
self.log = logger
|
||||
self.timeout_s = timeout_s
|
||||
@@ -301,7 +299,7 @@ class PollScheduler:
|
||||
self._poll_error_count += 1
|
||||
if self.log: self.log.warning("sweep cbor: %s", e)
|
||||
return
|
||||
indexed = self.sweep_index(tree)
|
||||
indexed = self.cache.index_device_tree(tree)
|
||||
for href, rep in indexed.items():
|
||||
with self._defer_lock:
|
||||
if self._defer_until.get(href, 0) > time.monotonic():
|
||||
@@ -58,6 +58,23 @@ class StateCache:
|
||||
merged.update(body)
|
||||
return self.apply_rep(href, merged, source='optimistic')
|
||||
|
||||
@staticmethod
|
||||
def index_device_tree(device0_body) -> dict[str, dict]:
|
||||
"""Turn a /device/0 CBOR list-of-{href, rep} sweep response into
|
||||
a dict keyed by href. Entry [0] is the device-level rep itself
|
||||
and isn't useful here, so it's skipped.
|
||||
|
||||
Replaces the old standalone sensors.index_links — folded in
|
||||
here because every current and future caller immediately feeds
|
||||
the result into apply_rep on this same cache."""
|
||||
out: dict[str, dict] = {}
|
||||
if not isinstance(device0_body, list):
|
||||
return out
|
||||
for entry in device0_body[1:]:
|
||||
if isinstance(entry, dict) and 'href' in entry:
|
||||
out[entry['href']] = entry.get('rep') or {}
|
||||
return out
|
||||
|
||||
def get(self, href: str) -> Optional[dict]:
|
||||
with self._lock:
|
||||
return self.links.get(href)
|
||||
@@ -24,14 +24,13 @@ import time
|
||||
import cbor2
|
||||
|
||||
from .appliances.base import ApplianceDescriptor, bridge_diagnostic_discovery
|
||||
from .coap_dtls import DtlsCoapSession, fmt_code
|
||||
from .config import ApplianceConfig, SharedConfig
|
||||
from .keepalive import KeepaliveTask
|
||||
from .logger import bridge_logger
|
||||
from .observe_refresh import ObserveRefreshTask
|
||||
from .poll_scheduler import PollScheduler
|
||||
from .sensors import index_links
|
||||
from .state_cache import StateCache
|
||||
from protocol.dtls_session import DtlsCoapSession, fmt_code
|
||||
from ocf.keepalive import KeepaliveTask
|
||||
from ocf.observe_refresh import ObserveRefreshTask
|
||||
from ocf.poll_scheduler import PollScheduler
|
||||
from ocf.state_cache import StateCache
|
||||
|
||||
|
||||
DEBUG_BRIDGE = os.environ.get('DEBUG_BRIDGE') == '1'
|
||||
@@ -295,7 +294,6 @@ class PushBridge:
|
||||
scheduler = PollScheduler(
|
||||
sess, self.cache,
|
||||
tiers=self.descriptor.poll_tiers,
|
||||
sweep_index_fn=index_links,
|
||||
is_active_fn=self.descriptor.is_active,
|
||||
logger=self.log,
|
||||
)
|
||||
@@ -357,7 +355,7 @@ class PushBridge:
|
||||
# During the seed we want the cache populated without triggering
|
||||
# a publish per resource — gate the on_change callback off until
|
||||
# the publish gate opens just below.
|
||||
for href, rep in index_links(body).items():
|
||||
for href, rep in StateCache.index_device_tree(body).items():
|
||||
if href not in self.cache.links:
|
||||
self.cache.apply_rep(href, rep, source='seed')
|
||||
self.last_seed_ts = time.time()
|
||||
|
||||
@@ -1,19 +0,0 @@
|
||||
"""Shared sensor helpers.
|
||||
|
||||
Appliance-specific flattening lives in samsung_appliance/appliances/*.py;
|
||||
this module only carries utilities that every descriptor uses (currently
|
||||
the /device/0 link-dict indexer).
|
||||
"""
|
||||
|
||||
|
||||
def index_links(device0_body):
|
||||
"""Turn the /device/0 CBOR list-of-{href, rep} into a dict keyed
|
||||
by href. The first list entry is the device-level rep itself and
|
||||
isn't useful here, so skip it."""
|
||||
out = {}
|
||||
if not isinstance(device0_body, list):
|
||||
return out
|
||||
for entry in device0_body[1:]:
|
||||
if isinstance(entry, dict) and 'href' in entry:
|
||||
out[entry['href']] = entry.get('rep') or {}
|
||||
return out
|
||||
@@ -0,0 +1,40 @@
|
||||
# tests/test_state_cache.py
|
||||
from ocf.state_cache import StateCache
|
||||
|
||||
|
||||
class _FakeDescriptor:
|
||||
on_observation = None
|
||||
|
||||
|
||||
def test_index_device_tree_skips_device_level_entry_at_index_zero():
|
||||
tree = [
|
||||
{'href': '/device/0', 'rep': {'n': 'Dryer'}}, # device-level, skipped
|
||||
{'href': '/mode/vs/0', 'rep': {'x': 1}},
|
||||
{'href': '/power/vs/0', 'rep': {'y': 2}},
|
||||
]
|
||||
indexed = StateCache.index_device_tree(tree)
|
||||
assert indexed == {'/mode/vs/0': {'x': 1}, '/power/vs/0': {'y': 2}}
|
||||
assert '/device/0' not in indexed
|
||||
|
||||
|
||||
def test_index_device_tree_stub_entry_becomes_empty_dict():
|
||||
tree = [
|
||||
{'href': '/device/0', 'rep': {}},
|
||||
{'href': '/oven/vs/0'}, # no 'rep' key at all — a stub resource
|
||||
]
|
||||
indexed = StateCache.index_device_tree(tree)
|
||||
assert indexed == {'/oven/vs/0': {}}
|
||||
|
||||
|
||||
def test_index_device_tree_non_list_input_returns_empty():
|
||||
assert StateCache.index_device_tree({'not': 'a list'}) == {}
|
||||
assert StateCache.index_device_tree(None) == {}
|
||||
|
||||
|
||||
def test_apply_rep_reports_change_and_updates_cache():
|
||||
cache = StateCache(_FakeDescriptor())
|
||||
changed = cache.apply_rep('/mode/vs/0', {'x': 1}, source='poll')
|
||||
assert changed is True
|
||||
assert cache.get('/mode/vs/0') == {'x': 1}
|
||||
unchanged = cache.apply_rep('/mode/vs/0', {'x': 1}, source='poll')
|
||||
assert unchanged is False
|
||||
Reference in New Issue
Block a user