Skip to content

Commit 1a47754

Browse files
authored
Merge pull request #71 from Moballo-LLC/codex/observe-lifecycle-followups
Fix probationary Observe delivery and quiesced cleanup
2 parents a0cd5a1 + eadb2f9 commit 1a47754

5 files changed

Lines changed: 272 additions & 31 deletions

File tree

README.md

Lines changed: 9 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -81,9 +81,11 @@ An RFC 7641 relation is confirmed only by a valid Observe response option;
8181
duplicate and stale 24-bit sequence values are not delivered. Some older
8282
Samsung firmware omits that option. For those devices, a plain initial `2.05`
8383
is probationary until a later packet arrives on the same token with a different
84-
Message ID. Optional `on_observe_pending`, `on_legacy_notification`, and
85-
`on_observe_error` constructor callbacks let consumers keep that compatibility
86-
path distinct from confirmed RFC notifications and ordinary polling.
84+
Message ID. Its complete representation is still delivered through
85+
`on_notification`; a blockwise representation is re-read before delivery.
86+
Optional `on_observe_pending`, `on_legacy_notification`, and `on_observe_error`
87+
constructor callbacks let consumers keep that compatibility path distinct from
88+
confirmed RFC notifications and ordinary polling.
8789

8890
Periodic renewal can target only the relations that need it; unrelated
8991
observations remain active. Existing query variants are preserved unless the
@@ -143,9 +145,10 @@ Interrupted attempts raise `SessionClosedError`.
143145
Hosts that stop network work before their blocking executor drains can use the
144146
session's two-phase shutdown. `quiesce_for_close()` is terminal: it interrupts
145147
an in-progress handshake, wakes pending requests and notification refetches,
146-
and rejects new work while retaining an established DTLS socket. A subsequent
147-
`close()` flushes the authenticated close-notify record before closing that
148-
socket:
148+
and rejects new work while retaining an established DTLS socket and active
149+
Observe relation metadata. A subsequent `close()` paces explicit Observe
150+
deregistrations, flushes the authenticated close-notify record, and then closes
151+
that socket:
149152

150153
```python
151154
sess.quiesce_for_close() # safe from the host's early shutdown phase

smartthings_local/protocol/dtls_session.py

Lines changed: 73 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -469,6 +469,11 @@ def __init__(self, host, port, cert_path=None, key_path=None, *,
469469
self._lifecycle_lock = threading.Lock()
470470

471471
self._send_lock = threading.Lock()
472+
# Only the thread performing orderly close may send an Observe
473+
# deregistration after terminal quiescence. A thread-local exception
474+
# keeps concurrent application senders blocked without changing the
475+
# existing private send-hook signatures used by test/session adapters.
476+
self._orderly_close_send_thread_id = None
472477
# Guards the MID/token counters and pending-request registries.
473478
# The refetch worker makes the session its own second concurrent
474479
# get() caller, so two threads can mint tokens at once; without
@@ -533,6 +538,18 @@ def pace(self) -> None:
533538
if remaining > 0:
534539
self._stop.wait(remaining)
535540

541+
def _pace_orderly_close(self) -> None:
542+
"""Honor request spacing after terminal quiescence.
543+
544+
``quiesce_for_close()`` sets ``_stop`` so ordinary paced work wakes
545+
immediately. A later orderly close still has to space its explicit
546+
Observe deregistrations, so this teardown-only path cannot wait on the
547+
already-set event.
548+
"""
549+
remaining = self._min_req_interval - (time.monotonic() - self._last_send_ts)
550+
if remaining > 0:
551+
time.sleep(remaining)
552+
536553
# ---- lifecycle ---------------------------------------------------
537554

538555
def connect(
@@ -712,8 +729,16 @@ def _send_observe_dereg(self, tok, path_segs, query=()):
712729
opts.append((URI_QUERY, value.encode()))
713730
opts.append((OBSERVE, OBSERVE_DEREGISTER))
714731
opts.append((ACCEPT, CF_CBOR))
715-
self._send_dgram(
716-
build_coap(TYPE_CON, METHOD_GET, mid, tok, opts))
732+
self._send_dgram(build_coap(TYPE_CON, METHOD_GET, mid, tok, opts))
733+
734+
def _send_observe_dereg_after_quiesce(self, tok, path_segs, query=()):
735+
"""Permit one orderly-close deregistration on the closing thread."""
736+
previous = self._orderly_close_send_thread_id
737+
self._orderly_close_send_thread_id = threading.get_ident()
738+
try:
739+
self._send_observe_dereg(tok, path_segs, query)
740+
finally:
741+
self._orderly_close_send_thread_id = previous
717742

718743
@staticmethod
719744
def _send_close_notify(connection, sock):
@@ -764,17 +789,25 @@ def _close_orderly(self):
764789
# Send dereg for every active observation while the conn is
765790
# still healthy. Tiny sleep lets the records reach the wire
766791
# before we shut DTLS down.
767-
if (not self._lifecycle_cancel.is_set() and self.conn is not None
768-
and self._observe_tokens):
792+
if self.conn is not None and self._observe_tokens:
793+
quiesced = self._lifecycle_cancel.is_set()
769794
with self._state_lock:
770795
observations = tuple(self._observe_tokens.items())
771796
observe_queries = dict(self._observe_queries)
772797
for tok, href in observations:
773798
segs = [s for s in href.split('/') if s]
774799
try:
775-
self.pace()
776-
self._send_observe_dereg(
777-
tok, segs, observe_queries.get(tok, ()))
800+
if quiesced:
801+
self._pace_orderly_close()
802+
self._send_observe_dereg_after_quiesce(
803+
tok,
804+
segs,
805+
observe_queries.get(tok, ()),
806+
)
807+
else:
808+
self.pace()
809+
self._send_observe_dereg(
810+
tok, segs, observe_queries.get(tok, ()))
778811
except Exception as e:
779812
logger.warning("dereg %s: %s", href, e)
780813
time.sleep(0.1)
@@ -918,7 +951,7 @@ def _clear_observe_relations(self):
918951
self._observe_sequences.clear()
919952

920953
def _observe_relation_active(self, href, query, legacy):
921-
"""Return whether one confirmed relation still owns this identity."""
954+
"""Return whether one relation still owns this callback identity."""
922955
with self._state_lock:
923956
for tok, observed_href in self._observe_tokens.items():
924957
if observed_href != href or \
@@ -927,7 +960,14 @@ def _observe_relation_active(self, href, query, legacy):
927960
if legacy:
928961
if tok in self._legacy_observe_tokens:
929962
return True
930-
elif tok in self._observe_sequences:
963+
elif tok in self._observe_sequences or (
964+
tok in self._observe_plain_response_mids
965+
and tok not in self._legacy_observe_tokens):
966+
# A probationary optionless response is not proof of an
967+
# Observe relation, but its complete representation still
968+
# belongs on the ordinary notification callback. Keep a
969+
# Block2 refetch alive until the token is retired or later
970+
# proves the legacy relation.
931971
return True
932972
return False
933973

@@ -955,9 +995,17 @@ def _observe_sequence_is_fresh(previous, current, received_at):
955995

956996
def _send_dgram(self, datagram):
957997
"""Send a CoAP datagram. Holds the send lock for the
958-
BIO-drain so two writers can't interleave records."""
998+
BIO-drain so two writers can't interleave records.
999+
1000+
The orderly-close deregistration helper grants only its calling thread
1001+
a teardown send after application workers have joined. All ordinary
1002+
request paths remain blocked once terminal quiescence begins.
1003+
"""
9591004
with self._send_lock:
960-
if self._lifecycle_cancel.is_set() or self.conn is None:
1005+
orderly_close_send = (
1006+
self._orderly_close_send_thread_id == threading.get_ident())
1007+
if (self._lifecycle_cancel.is_set() and not orderly_close_send) or \
1008+
self.conn is None:
9611009
raise SessionClosedError()
9621010
send_failed = False
9631011
try:
@@ -1058,7 +1106,12 @@ def _reader_loop(self):
10581106
with self._refetch_cond:
10591107
self._refetch_pending.clear()
10601108
self._refetch_cond.notify_all()
1061-
self._clear_observe_relations()
1109+
# Two-phase shutdown retains relation metadata for close(), which
1110+
# runs after application workers have joined and sends the paced
1111+
# deregistration sweep. Unexpected reader death has no later
1112+
# orderly phase and must still retire everything immediately.
1113+
if not self._lifecycle_cancel.is_set():
1114+
self._clear_observe_relations()
10621115

10631116
def _dispatch_coap(self, datagram):
10641117
try:
@@ -1221,7 +1274,14 @@ def _dispatch_coap(self, datagram):
12211274
except Exception as e:
12221275
logger.debug(
12231276
"observe pending callback %s: %s", href, e)
1224-
return
1277+
with self._state_lock:
1278+
if self._observe_tokens.get(tok) != href or \
1279+
self._observe_queries.get(tok, ()) != \
1280+
observe_query or \
1281+
self._observe_plain_response_mids.get(tok) != \
1282+
mid or tok in self._legacy_observe_tokens or \
1283+
tok in self._observe_sequences:
1284+
return
12251285
# RFC 7959 §2.6: a notification carries only the first block
12261286
# of the representation. Handing the callback a partial CBOR
12271287
# buffer is what #39 was about, so anything with M=1 (or a

tests/test_observe_operations.py

Lines changed: 50 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -7,7 +7,7 @@
77

88
import pytest
99

10-
from smartthings_local.errors import EndpointError
10+
from smartthings_local.errors import EndpointError, SessionClosedError
1111
from smartthings_local.protocol.coap import (
1212
OBSERVE,
1313
URI_PATH,
@@ -345,18 +345,64 @@ def test_normal_close_paces_every_exact_deregister_and_clears_state(
345345
assert session._observe_sequences == {}
346346

347347

348-
def test_quiesced_close_skips_deregister_pacing():
348+
def test_quiesced_close_paces_every_deregister_for_orderly_teardown():
349349
session = _session()
350350
session.sock = _Socket()
351-
_add_relation(session, ["mode", "vs", "0"])
351+
token = _add_relation(session, ["mode", "vs", "0"])
352352
session.pace = Mock()
353+
session._pace_orderly_close = Mock()
353354
session._send_observe_dereg = Mock()
354355

355356
session.quiesce_for_close()
356357
session.close()
357358

358359
session.pace.assert_not_called()
359-
session._send_observe_dereg.assert_not_called()
360+
session._pace_orderly_close.assert_called_once_with()
361+
session._send_observe_dereg.assert_called_once_with(
362+
token,
363+
["mode", "vs", "0"],
364+
(),
365+
)
366+
367+
368+
def test_quiesced_close_preserves_existing_send_override_signature():
369+
session = _session()
370+
session.sock = _Socket()
371+
token = _add_relation(session, ["mode", "vs", "0"])
372+
sent = []
373+
session._send_dgram = sent.append
374+
session._pace_orderly_close = Mock()
375+
376+
session.quiesce_for_close()
377+
session.close()
378+
379+
assert len(sent) == 1
380+
assert parse_coap(sent[0])[3] == token
381+
382+
383+
def test_orderly_close_send_permission_is_thread_local():
384+
session = _session()
385+
session.quiesce_for_close()
386+
outcomes = []
387+
388+
def try_application_send():
389+
try:
390+
session._send_dgram(b"application request")
391+
except Exception as error: # noqa: BLE001 - captured for assertion
392+
outcomes.append(error)
393+
394+
def during_teardown_send(_token, _path, _query):
395+
thread = threading.Thread(target=try_application_send)
396+
thread.start()
397+
thread.join()
398+
399+
session._send_observe_dereg = during_teardown_send
400+
session._send_observe_dereg_after_quiesce(
401+
b"o", ["mode", "vs", "0"], ()
402+
)
403+
404+
assert len(outcomes) == 1
405+
assert isinstance(outcomes[0], SessionClosedError)
360406

361407

362408
def test_reader_exit_clears_all_relation_state():

tests/test_observe_relations.py

Lines changed: 88 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -12,7 +12,6 @@
1212
BLOCK2,
1313
OBSERVE,
1414
TYPE_ACK,
15-
TYPE_CON,
1615
TYPE_NON,
1716
URI_QUERY,
1817
block_value,
@@ -71,7 +70,7 @@ def test_subscribe_registers_path_and_query_before_immediate_response():
7170
def send(datagram):
7271
request = parse_coap(datagram)
7372
requests.append(request)
74-
_mtype, _code, mid, token, options, _payload = request
73+
_mtype, _code, mid, token, _options, _payload = request
7574
assert session._observe_tokens[token] == "/mode/vs/0"
7675
assert session._observe_queries[token] == (
7776
"if=oic.if.a",
@@ -122,15 +121,20 @@ def test_plain_initial_response_needs_different_mid_to_confirm_legacy():
122121
_notify(session, token, 10, b"initial-retransmit")
123122

124123
assert pending == ["/doors/vs/0"]
125-
assert delivered == []
124+
assert delivered == [
125+
("standard", "/doors/vs/0", b"initial"),
126+
]
127+
assert session._observe_plain_response_mids[token] == 10
126128
assert token not in session._legacy_observe_tokens
129+
assert token not in session._observe_sequences
127130

128131
_notify(session, token, 11, b"changed")
129132
_notify(session, token, 11, b"changed-retransmit")
130133
_notify(session, token, 12, b"changed-again")
131134

132135
assert token in session._legacy_observe_tokens
133136
assert delivered == [
137+
("standard", "/doors/vs/0", b"initial"),
134138
("legacy", "/doors/vs/0", b"changed"),
135139
("legacy", "/doors/vs/0", b"changed-again"),
136140
]
@@ -147,7 +151,87 @@ def test_legacy_notification_falls_back_to_main_callback():
147151
_notify(session, token, 20, b"initial")
148152
_notify(session, token, 21, b"changed")
149153

150-
assert delivered == [("/doors/vs/0", b"changed")]
154+
assert delivered == [
155+
("/doors/vs/0", b"initial"),
156+
("/doors/vs/0", b"changed"),
157+
]
158+
159+
160+
def test_probationary_blockwise_initial_queues_refetch_without_partial_delivery():
161+
delivered = []
162+
pending = []
163+
session = _session(
164+
on_notification=lambda href, payload: delivered.append((href, payload)),
165+
on_observe_pending=pending.append,
166+
)
167+
session._send_dgram = Mock()
168+
session._queue_refetch = Mock()
169+
token = session.subscribe(
170+
["mode", "vs", "0"], query=("if=oic.if.a",)
171+
)
172+
173+
_notify(
174+
session,
175+
token,
176+
10,
177+
b"partial",
178+
options=((BLOCK2, block_value(0, 1, 6)),),
179+
)
180+
181+
assert pending == ["/mode/vs/0"]
182+
assert delivered == []
183+
session._queue_refetch.assert_called_once_with(
184+
"/mode/vs/0",
185+
("if=oic.if.a",),
186+
legacy=False,
187+
)
188+
189+
190+
def test_pending_callback_can_retire_relation_before_initial_delivery():
191+
delivered = []
192+
holder = {}
193+
194+
def retire(_href):
195+
with holder["session"]._state_lock:
196+
holder["session"]._retire_observe_token_locked(holder["token"])
197+
198+
session = _session(
199+
on_notification=lambda href, payload: delivered.append((href, payload)),
200+
on_observe_pending=retire,
201+
)
202+
holder["session"] = session
203+
session._send_dgram = Mock()
204+
holder["token"] = session.subscribe(["mode", "vs", "0"])
205+
206+
_notify(session, holder["token"], 10, b"initial")
207+
208+
assert delivered == []
209+
assert holder["token"] not in session._observe_tokens
210+
211+
212+
def test_probationary_blockwise_representation_is_refetched_before_delivery():
213+
delivered = []
214+
session = _session(
215+
on_notification=lambda href, payload: delivered.append((href, payload))
216+
)
217+
session._blockwise_get = Mock(
218+
return_value=(0x45, b"complete", 2, b"fresh")
219+
)
220+
token = b"\x41"
221+
href = "/mode/vs/0"
222+
query = ("if=oic.if.a",)
223+
session._observe_tokens[token] = href
224+
session._observe_queries[token] = query
225+
session._observe_plain_response_mids[token] = 10
226+
227+
session._refetch_one((href, query, False), 1)
228+
229+
session._blockwise_get.assert_called_once_with(
230+
["mode", "vs", "0"],
231+
query,
232+
dtls_session._REFETCH_TIMEOUT_S,
233+
)
234+
assert delivered == [(href, b"complete")]
151235

152236

153237
def test_rejected_observe_retires_every_relation_index():

0 commit comments

Comments
 (0)