44"""The controller's side of a calibration session: one at a time, with events.
55
66The capture loop itself is `CalibrationSession`; this drives it from async
7- routes. It owns the swarmit client, built on first capture rather than at
8- start so a controller with no fleet in reach still serves the routes ; it
7+ routes. It owns the swarmit client, built when a session starts and rebuilt
8+ when the robot changes, and a start with no fleet in reach still opens ; it
99turns every state change into exactly one WebSocket notification, in the
1010order the changes happened; and it routes a capture that arrives from the
11- robot's own trigger to the outstanding point.
11+ robot's own button to the outstanding point.
1212
1313The blocking capture runs in a worker thread. Progress is published from
1414there by handing the coroutine back to the loop and waiting for it, so the
2222
2323from dotbot .calibration .ota import CAPTURE_READS_DEFAULT
2424from dotbot .calibration .points import resolve_placement_points
25+ from dotbot .calibration .push import PushRefused , gate_push , push_worklist
2526from dotbot .calibration .session import (
2627 CalibrationSession ,
2728 SessionError ,
3031from dotbot .logger import LOGGER
3132from dotbot .site import Site
3233
34+ # Seconds between two checks for a button press given up incomplete.
35+ EXPIRED_PRESS_POLL_INTERVAL = 1.0
36+
3337# The swarmit log-event tag a raw-count capture carries. Imported lazily so
3438# the swarmit protocol registry stays out of PyDotBot test collection.
3539_CAPTURE_TAG : int | None = None
@@ -47,22 +51,19 @@ def capture_tag() -> int:
4751class SessionDriver :
4852 """One calibration session at a time, with its transport and its events.
4953
50- `client_factory` takes the device address and returns a swarmit client;
51- `stale_devices` reports which robots do not hold the calibration in use.
52- Both are injected so a test drives the whole loop without a fleet.
54+ `client_factory` takes the device address and returns a swarmit client,
55+ injected so a test drives the whole loop without a fleet.
5356 """
5457
5558 def __init__ (
5659 self ,
5760 client_factory : Callable [[str ], Any ],
5861 notify : Callable [[dict | None ], Any ],
5962 site : Site | None = None ,
60- stale_devices : Callable [[], list [str ]] | None = None ,
6163 stream_factory : Callable [[Any , str , Callable ], Any ] | None = None ,
6264 ):
6365 self ._client_factory = client_factory
6466 self ._notify = notify
65- self ._stale_devices = stale_devices or (lambda : [])
6667 self ._stream_factory = stream_factory or _default_stream
6768 self .site = site or Site ()
6869 self .session : CalibrationSession | None = None
@@ -72,6 +73,7 @@ def __init__(
7273 self ._device = ""
7374 self ._lock = asyncio .Lock ()
7475 self ._loop : asyncio .AbstractEventLoop | None = None
76+ self ._expiry_watch : asyncio .Task | None = None
7577
7678 # -- state
7779
@@ -114,6 +116,12 @@ async def start(
114116 session .device = device .upper ()
115117 session .area = area
116118 self .session = session
119+ self ._loop = asyncio .get_running_loop ()
120+ # Button captures come unrequested from any robot, so the stream
121+ # is listening from the start rather than from the first capture.
122+ await asyncio .to_thread (self ._listen_for_buttons , session .device )
123+ if self ._expiry_watch is None or self ._expiry_watch .done ():
124+ self ._expiry_watch = asyncio .create_task (self ._watch_expired_presses ())
117125 await self ._emit ()
118126 return session .as_dict ()
119127
@@ -166,18 +174,38 @@ async def save(self, tag: str = "") -> dict:
166174 "session" : session .as_dict (),
167175 }
168176
169- async def push (self ) -> dict :
170- """Send the saved calibration and report which robots are still stale."""
177+ async def push (
178+ self , site_changed : bool = False , devices : list [str ] | None = None
179+ ) -> dict :
180+ """Check `devices`, send them the saved calibration, report who is still stale.
181+
182+ No devices, None or empty, is the whole swarm, as for flash and start.
183+ """
184+ targets = sorted ({d .upper () for d in devices }) if devices else None
171185 async with self ._lock :
172186 session = self ._require ()
173187 payload = session .push_payload ()
174- client = await asyncio .to_thread (self ._ensure_client , session .device )
175- await asyncio .to_thread (client .send_lh2_calibration , payload )
188+ client = await asyncio .to_thread (self ._ensure_client , "" )
189+ try :
190+ try :
191+ check = await asyncio .to_thread (
192+ gate_push , client , session .saved , site_changed , targets
193+ )
194+ except PushRefused as exc :
195+ raise SessionError (f"push refused: { exc } " ) from exc
196+ await asyncio .to_thread (
197+ client .send_lh2_calibration , payload , check .send_to
198+ )
199+ stale = await asyncio .to_thread (
200+ push_worklist , client , session .saved , check .addresses
201+ )
202+ finally :
203+ await asyncio .to_thread (self ._listen_for_buttons , "" )
176204 await self ._emit ()
177205 return {
178206 "id" : session .saved_id ,
179207 "bytes" : len (payload ),
180- "stale" : self . _stale_devices () ,
208+ "stale" : stale ,
181209 }
182210
183211 async def abandon (self ) -> dict :
@@ -190,27 +218,73 @@ async def abandon(self) -> dict:
190218
191219 # -- the robot's own trigger
192220
193- def on_idle_records (self , records : list ) -> None :
194- """A capture that arrived without a request: point k, or dropped."""
195- session = self .session
196- if session is None :
197- self .logger .info (
198- "LH2 capture arrived with no calibration session open; dropped" ,
199- records = len (records ),
200- )
201- return
202- point = session .store_records (records )
203- if point is None :
221+ def on_button_capture (self , capture : Any ) -> Any :
222+ """A capture from the robot's own button, from the stream's reader thread.
223+
224+ Stored on the event loop under the session lock; returns the
225+ concurrent future of that, or None when no session ever started.
226+ """
227+ if self ._loop is None :
204228 self .logger .info (
205- "LH2 capture arrived with every point captured ; dropped" ,
206- records = len ( records ) ,
229+ "LH2 button capture arrived with no calibration session open ; dropped" ,
230+ device = capture . device ,
207231 )
208- return
209- self . logger . info (
210- "LH2 capture stored from the robot's own trigger" , point = point . index
232+ return None
233+ return asyncio . run_coroutine_threadsafe (
234+ self . _store_button_capture ( capture ), self . _loop
211235 )
212- if self ._loop is not None :
213- asyncio .run_coroutine_threadsafe (self ._emit (), self ._loop )
236+
237+ async def _watch_expired_presses (self ) -> None :
238+ """Report each button press given up incomplete, until the session ends."""
239+ while self .session is not None :
240+ await asyncio .sleep (EXPIRED_PRESS_POLL_INTERVAL )
241+ async with self ._lock :
242+ if self .session is None or self ._stream is None :
243+ continue
244+ expired = self ._stream .expired_presses ()
245+ if not expired :
246+ continue
247+ self .session .error = "; " .join (
248+ f"incomplete capture from { addr } (press { press } ): a chunk "
249+ "never arrived; press again"
250+ for addr , press in expired
251+ )
252+ await self ._emit ()
253+
254+ async def _store_button_capture (self , capture : Any ) -> None :
255+ """Point k, or dropped."""
256+ async with self ._lock :
257+ session = self .session
258+ if session is None :
259+ self .logger .info (
260+ "LH2 button capture arrived with no calibration session open; dropped" ,
261+ device = capture .device ,
262+ )
263+ return
264+ if capture .lost :
265+ self .logger .warning (
266+ "LH2 button captures lost" , device = capture .device , lost = capture .lost
267+ )
268+ try :
269+ point = session .store_reads (capture .reads , capture .device )
270+ except SessionError as exc :
271+ session .error = str (exc )
272+ else :
273+ if point is None :
274+ self .logger .info (
275+ "LH2 button capture arrived with every point captured; dropped" ,
276+ device = capture .device ,
277+ )
278+ return
279+ session .error = ""
280+ # The robot that pressed is the one capturing, as if chosen.
281+ session .device = point .device
282+ self .logger .info (
283+ "LH2 capture stored from the robot's own button" ,
284+ device = capture .device ,
285+ point = point .index ,
286+ )
287+ await self ._emit ()
214288
215289 # -- transport
216290
@@ -225,6 +299,9 @@ def _require(self) -> CalibrationSession:
225299 def _ensure_client (self , device : str ) -> Any :
226300 if self ._client is None or device != self ._device :
227301 self ._close_stream ()
302+ if self ._client is not None :
303+ self ._client .__exit__ (None , None , None )
304+ self ._client = None
228305 self ._client = self ._client_factory (device )
229306 self ._client .__enter__ ()
230307 self ._device = device
@@ -233,10 +310,19 @@ def _ensure_client(self, device: str) -> Any:
233310 def _ensure_stream (self , device : str ) -> Any :
234311 client = self ._ensure_client (device )
235312 if self ._stream is None :
236- self ._stream = self ._stream_factory (client , device , self .on_idle_records )
313+ self ._stream = self ._stream_factory (client , device , self .on_button_capture )
237314 self ._stream .__enter__ ()
238315 return self ._stream
239316
317+ def _listen_for_buttons (self , device : str ) -> None :
318+ try :
319+ self ._ensure_stream (device )
320+ except Exception as exc : # a fleet out of reach is not a failure
321+ self .logger .warning (
322+ "Not listening for button captures until a capture reaches the fleet" ,
323+ error = str (exc ),
324+ )
325+
240326 def _close_stream (self ) -> None :
241327 if self ._stream is not None :
242328 self ._stream .__exit__ (None , None , None )
@@ -247,9 +333,9 @@ def bind_loop(self, loop: asyncio.AbstractEventLoop) -> None:
247333 self ._loop = loop
248334
249335
250- def _default_stream (client : Any , device : str , on_idle_records : Callable ) -> Any :
336+ def _default_stream (client : Any , device : str , on_button_capture : Callable ) -> Any :
251337 from dotbot .calibration .ota import CaptureSession
252338
253339 return CaptureSession (
254- client , device , capture_tag (), on_idle_records = on_idle_records
340+ client , device , capture_tag (), on_button_capture = on_button_capture
255341 )
0 commit comments