aranet_core/
history.rs

1//! Historical data download.
2//!
3//! This module provides functionality to download historical sensor
4//! readings stored on an Aranet device.
5//!
6//! # Supported Devices
7//!
8//! | Device | History Support | Notes |
9//! |--------|-----------------|-------|
10//! | Aranet4 | Full | CO₂, temperature, pressure, humidity |
11//! | Aranet2 | Full | Temperature, humidity |
12//! | AranetRn+ (Radon) | Full | Radon, temperature, pressure, humidity |
13//! | Aranet Radiation | Not supported | BLE protocol undocumented by manufacturer |
14//!
15//! **Note:** Aranet Radiation devices do not support history download because
16//! the BLE protocol for historical radiation data is not publicly documented
17//! by SAF Tehnika. Attempting to download history will return [`Error::Unsupported`].
18//! Use [`Device::read_current()`](crate::device::Device::read_current) for
19//! current radiation readings.
20//!
21//! # Index Convention
22//!
23//! **All history indices are 1-based**, following the Aranet device protocol:
24//! - Index 1 = oldest reading
25//! - Index N = newest reading (where N = total_readings)
26//!
27//! This matches the device's internal indexing. When specifying ranges:
28//! ```ignore
29//! let options = HistoryOptions {
30//!     start_index: Some(1),    // First reading
31//!     end_index: Some(100),    // 100th reading
32//!     ..Default::default()
33//! };
34//! ```
35//!
36//! # Protocols
37//!
38//! Aranet devices support two history protocols:
39//! - **V1**: Notification-based (older devices) - uses characteristic notifications
40//! - **V2**: Read-based (newer devices, preferred) - direct read/write operations
41
42use std::collections::BTreeMap;
43use std::sync::Arc;
44use std::time::Duration;
45
46use bytes::Buf;
47use time::OffsetDateTime;
48use tokio::time::sleep;
49use tracing::{debug, info, warn};
50
51use crate::commands::{HISTORY_V1_REQUEST, HISTORY_V2_REQUEST};
52use crate::device::Device;
53use crate::error::{Error, Result};
54use crate::uuid::{COMMAND, HISTORY_V2, READ_INTERVAL, SECONDS_SINCE_UPDATE, TOTAL_READINGS};
55use aranet_types::HistoryRecord;
56
57/// Progress information for history download.
58#[derive(Debug, Clone)]
59pub struct HistoryProgress {
60    /// Current parameter being downloaded.
61    pub current_param: HistoryParam,
62    /// Parameter index (1-based, e.g., 1 of 4).
63    pub param_index: usize,
64    /// Total number of parameters to download.
65    pub total_params: usize,
66    /// Number of values downloaded for current parameter.
67    pub values_downloaded: usize,
68    /// Total values to download for current parameter.
69    pub total_values: usize,
70    /// Overall progress (0.0 to 1.0).
71    pub overall_progress: f32,
72}
73
74impl HistoryProgress {
75    /// Create a new progress struct.
76    pub fn new(
77        param: HistoryParam,
78        param_idx: usize,
79        total_params: usize,
80        total_values: usize,
81    ) -> Self {
82        Self {
83            current_param: param,
84            param_index: param_idx,
85            total_params,
86            values_downloaded: 0,
87            total_values,
88            overall_progress: 0.0,
89        }
90    }
91
92    fn update(&mut self, values_downloaded: usize) {
93        self.values_downloaded = values_downloaded;
94        let param_progress = if self.total_values > 0 {
95            values_downloaded as f32 / self.total_values as f32
96        } else {
97            1.0
98        };
99        // Guard against division by zero when total_params is 0
100        if self.total_params == 0 {
101            self.overall_progress = 1.0;
102            return;
103        }
104        let base_progress = (self.param_index - 1) as f32 / self.total_params as f32;
105        let param_contribution = param_progress / self.total_params as f32;
106        self.overall_progress = base_progress + param_contribution;
107    }
108}
109
110/// Type alias for progress callback function.
111pub type ProgressCallback = Arc<dyn Fn(HistoryProgress) + Send + Sync>;
112
113/// Type alias for checkpoint callback function.
114pub type CheckpointCallback = Arc<dyn Fn(HistoryCheckpoint) + Send + Sync>;
115
116/// Checkpoint data for resuming interrupted history downloads.
117///
118/// This can be serialized and saved to disk to allow resuming downloads
119/// after disconnection or application restart.
120#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
121pub struct HistoryCheckpoint {
122    /// Device identifier this checkpoint belongs to.
123    pub device_id: String,
124    /// The parameter currently being downloaded.
125    pub current_param: HistoryParamCheckpoint,
126    /// Index where download should resume for current parameter.
127    pub resume_index: u16,
128    /// Total readings on the device when checkpoint was created.
129    pub total_readings: u16,
130    /// Which parameters have been fully downloaded.
131    pub completed_params: Vec<HistoryParamCheckpoint>,
132    /// Timestamp when checkpoint was created.
133    pub created_at: time::OffsetDateTime,
134    /// Downloaded values for completed parameters (serialized).
135    pub downloaded_data: Option<PartialHistoryData>,
136}
137
138/// Serializable version of HistoryParam for checkpoints.
139#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
140pub enum HistoryParamCheckpoint {
141    Temperature,
142    Humidity,
143    Pressure,
144    Co2,
145    Humidity2,
146    Radon,
147}
148
149impl From<HistoryParam> for HistoryParamCheckpoint {
150    fn from(param: HistoryParam) -> Self {
151        match param {
152            HistoryParam::Temperature => HistoryParamCheckpoint::Temperature,
153            HistoryParam::Humidity => HistoryParamCheckpoint::Humidity,
154            HistoryParam::Pressure => HistoryParamCheckpoint::Pressure,
155            HistoryParam::Co2 => HistoryParamCheckpoint::Co2,
156            HistoryParam::Humidity2 => HistoryParamCheckpoint::Humidity2,
157            HistoryParam::Radon => HistoryParamCheckpoint::Radon,
158        }
159    }
160}
161
162impl From<HistoryParamCheckpoint> for HistoryParam {
163    fn from(param: HistoryParamCheckpoint) -> Self {
164        match param {
165            HistoryParamCheckpoint::Temperature => HistoryParam::Temperature,
166            HistoryParamCheckpoint::Humidity => HistoryParam::Humidity,
167            HistoryParamCheckpoint::Pressure => HistoryParam::Pressure,
168            HistoryParamCheckpoint::Co2 => HistoryParam::Co2,
169            HistoryParamCheckpoint::Humidity2 => HistoryParam::Humidity2,
170            HistoryParamCheckpoint::Radon => HistoryParam::Radon,
171        }
172    }
173}
174
175#[derive(Debug, Clone, Copy)]
176struct U16HistoryStep {
177    param: HistoryParam,
178    step: usize,
179    total_steps: usize,
180    next_param: Option<HistoryParamCheckpoint>,
181}
182
183/// Partially downloaded history data for checkpoint resume.
184#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
185pub struct PartialHistoryData {
186    pub co2_values: Vec<u16>,
187    pub temp_values: Vec<u16>,
188    pub pressure_values: Vec<u16>,
189    pub humidity_values: Vec<u16>,
190    pub radon_values: Vec<u32>,
191}
192
193impl HistoryCheckpoint {
194    /// Create a new checkpoint for starting a fresh download.
195    pub fn new(device_id: &str, total_readings: u16, first_param: HistoryParam) -> Self {
196        Self {
197            device_id: device_id.to_string(),
198            current_param: first_param.into(),
199            resume_index: 1,
200            total_readings,
201            completed_params: Vec::new(),
202            created_at: time::OffsetDateTime::now_utc(),
203            downloaded_data: Some(PartialHistoryData::default()),
204        }
205    }
206
207    /// Check if this checkpoint is still valid for the given device state.
208    pub fn is_valid(&self, current_total_readings: u16) -> bool {
209        // Checkpoint is valid if the device hasn't collected more readings
210        // (which would shift the indices)
211        self.total_readings == current_total_readings
212    }
213
214    /// Update the checkpoint after completing a parameter.
215    pub fn complete_param(&mut self, param: HistoryParam, values: Vec<u16>) {
216        self.completed_params.push(param.into());
217        if let Some(ref mut data) = self.downloaded_data {
218            match param {
219                HistoryParam::Co2 => data.co2_values = values,
220                HistoryParam::Temperature => data.temp_values = values,
221                HistoryParam::Pressure => data.pressure_values = values,
222                HistoryParam::Humidity | HistoryParam::Humidity2 => data.humidity_values = values,
223                HistoryParam::Radon => {} // Radon uses u32, handled separately
224            }
225        }
226    }
227
228    /// Update the checkpoint after completing a radon parameter.
229    pub fn complete_radon_param(&mut self, values: Vec<u32>) {
230        self.completed_params.push(HistoryParamCheckpoint::Radon);
231        if let Some(ref mut data) = self.downloaded_data {
232            data.radon_values = values;
233        }
234    }
235}
236
237/// Parameter types for history requests.
238#[derive(Debug, Clone, Copy, PartialEq, Eq)]
239#[repr(u8)]
240pub enum HistoryParam {
241    Temperature = 1,
242    Humidity = 2,
243    Pressure = 3,
244    Co2 = 4,
245    /// Humidity for Aranet2/Radon (different encoding).
246    Humidity2 = 5,
247    /// Radon concentration (Bq/m³) for AranetRn+.
248    Radon = 10,
249}
250
251/// Options for downloading history.
252///
253/// # Index Convention
254///
255/// Indices are **1-based** to match the Aranet device protocol:
256/// - `start_index: Some(1)` means the first (oldest) reading
257/// - `end_index: Some(100)` means the 100th reading
258/// - `start_index: None` defaults to 1 (beginning)
259/// - `end_index: None` defaults to total_readings (end)
260///
261/// # Progress Reporting
262///
263/// Use `with_progress` to receive updates during download:
264/// ```ignore
265/// let options = HistoryOptions::default()
266///     .with_progress(|p| println!("Progress: {:.1}%", p.overall_progress * 100.0));
267/// ```
268///
269/// # Adaptive Read Delay
270///
271/// Use `adaptive_delay` to automatically adjust delay based on signal quality:
272/// ```ignore
273/// let options = HistoryOptions::default().adaptive_delay(true);
274/// ```
275///
276/// # Resume Support
277///
278/// For long downloads, use checkpointing to allow resume on failure:
279/// ```ignore
280/// let checkpoint = HistoryCheckpoint::load("device_123")?;
281/// let options = HistoryOptions::default().resume_from(checkpoint);
282/// ```
283#[derive(Clone)]
284pub struct HistoryOptions {
285    /// Starting index (1-based, inclusive). If None, downloads from the beginning (index 1).
286    pub start_index: Option<u16>,
287    /// Ending index (1-based, inclusive). If None, downloads to the end (index = total_readings).
288    pub end_index: Option<u16>,
289    /// Delay between read operations to avoid overwhelming the device.
290    pub read_delay: Duration,
291    /// Progress callback (optional).
292    pub progress_callback: Option<ProgressCallback>,
293    /// Whether to use adaptive delay based on signal quality.
294    pub use_adaptive_delay: bool,
295    /// Checkpoint callback for saving progress during download (optional).
296    /// Called periodically with the current checkpoint state.
297    pub checkpoint_callback: Option<CheckpointCallback>,
298    /// How often to call the checkpoint callback (in records).
299    pub checkpoint_interval: usize,
300}
301
302impl std::fmt::Debug for HistoryOptions {
303    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
304        f.debug_struct("HistoryOptions")
305            .field("start_index", &self.start_index)
306            .field("end_index", &self.end_index)
307            .field("read_delay", &self.read_delay)
308            .field("progress_callback", &self.progress_callback.is_some())
309            .field("use_adaptive_delay", &self.use_adaptive_delay)
310            .field("checkpoint_callback", &self.checkpoint_callback.is_some())
311            .field("checkpoint_interval", &self.checkpoint_interval)
312            .finish()
313    }
314}
315
316impl Default for HistoryOptions {
317    fn default() -> Self {
318        Self {
319            start_index: None,
320            end_index: None,
321            read_delay: Duration::from_millis(50),
322            progress_callback: None,
323            use_adaptive_delay: false,
324            checkpoint_callback: None,
325            checkpoint_interval: 100, // Checkpoint every 100 records
326        }
327    }
328}
329
330impl HistoryOptions {
331    /// Create new history options with default settings.
332    #[must_use]
333    pub fn new() -> Self {
334        Self::default()
335    }
336
337    /// Set the starting index (1-based).
338    #[must_use]
339    pub fn start_index(mut self, index: u16) -> Self {
340        self.start_index = Some(index);
341        self
342    }
343
344    /// Set the ending index (1-based).
345    #[must_use]
346    pub fn end_index(mut self, index: u16) -> Self {
347        self.end_index = Some(index);
348        self
349    }
350
351    /// Set the delay between read operations.
352    #[must_use]
353    pub fn read_delay(mut self, delay: Duration) -> Self {
354        self.read_delay = delay;
355        self
356    }
357
358    /// Set a progress callback.
359    #[must_use]
360    pub fn with_progress<F>(mut self, callback: F) -> Self
361    where
362        F: Fn(HistoryProgress) + Send + Sync + 'static,
363    {
364        self.progress_callback = Some(Arc::new(callback));
365        self
366    }
367
368    /// Report progress if a callback is set.
369    pub fn report_progress(&self, progress: &HistoryProgress) {
370        if let Some(cb) = &self.progress_callback {
371            cb(progress.clone());
372        }
373    }
374
375    /// Enable or disable adaptive delay based on signal quality.
376    ///
377    /// When enabled, the read delay will be automatically adjusted based on
378    /// the connection's signal strength:
379    /// - Excellent signal: 30ms delay
380    /// - Good signal: 50ms delay
381    /// - Fair signal: 100ms delay
382    /// - Poor signal: 200ms delay
383    #[must_use]
384    pub fn adaptive_delay(mut self, enable: bool) -> Self {
385        self.use_adaptive_delay = enable;
386        self
387    }
388
389    /// Set a checkpoint callback for saving download progress.
390    ///
391    /// The callback will be invoked periodically (based on `checkpoint_interval`)
392    /// with the current checkpoint state, allowing recovery from interruptions.
393    #[must_use]
394    pub fn with_checkpoint<F>(mut self, callback: F) -> Self
395    where
396        F: Fn(HistoryCheckpoint) + Send + Sync + 'static,
397    {
398        self.checkpoint_callback = Some(Arc::new(callback));
399        self
400    }
401
402    /// Set how often to call the checkpoint callback (in records).
403    ///
404    /// Default: 100 records
405    #[must_use]
406    pub fn checkpoint_interval(mut self, interval: usize) -> Self {
407        self.checkpoint_interval = interval;
408        self
409    }
410
411    /// Resume from a previous checkpoint.
412    ///
413    /// This sets the start_index based on the checkpoint's resume position.
414    #[must_use]
415    pub fn resume_from(mut self, checkpoint: &HistoryCheckpoint) -> Self {
416        self.start_index = Some(checkpoint.resume_index);
417        self
418    }
419
420    /// Report a checkpoint if a callback is set.
421    pub fn report_checkpoint(&self, checkpoint: &HistoryCheckpoint) {
422        if let Some(cb) = &self.checkpoint_callback {
423            cb(checkpoint.clone());
424        }
425    }
426
427    /// Get the effective read delay, optionally adjusted for signal quality.
428    pub fn effective_read_delay(
429        &self,
430        signal_quality: Option<crate::device::SignalQuality>,
431    ) -> Duration {
432        if self.use_adaptive_delay
433            && let Some(quality) = signal_quality
434        {
435            return quality.recommended_read_delay();
436        }
437        self.read_delay
438    }
439}
440
441/// Information about the device's stored history.
442#[derive(Debug, Clone)]
443pub struct HistoryInfo {
444    /// Total number of readings stored.
445    pub total_readings: u16,
446    /// Measurement interval in seconds.
447    pub interval_seconds: u16,
448    /// Seconds since the last reading.
449    pub seconds_since_update: u16,
450}
451
452/// Wall-clock reference for converting device history indices to timestamps.
453///
454/// Captured once, immediately after `seconds_since_update` is read, so the
455/// timestamps don't shift by however long the download takes. Index
456/// `total_readings` is the newest record on the device; each lower index is
457/// one interval older.
458#[derive(Debug, Clone, Copy)]
459struct HistoryAnchor {
460    newest_record_at: OffsetDateTime,
461    total_readings: u16,
462    interval_seconds: u16,
463}
464
465impl HistoryAnchor {
466    fn new(info: &HistoryInfo, read_at: OffsetDateTime) -> Self {
467        Self {
468            newest_record_at: read_at
469                - time::Duration::seconds(i64::from(info.seconds_since_update)),
470            total_readings: info.total_readings,
471            interval_seconds: info.interval_seconds,
472        }
473    }
474
475    /// Timestamp of the record at 1-based device `index`.
476    fn timestamp_for(&self, index: u32) -> OffsetDateTime {
477        let readings_ago = i64::from(self.total_readings) - i64::from(index);
478        self.newest_record_at
479            - time::Duration::seconds(readings_ago * i64::from(self.interval_seconds))
480    }
481}
482
483impl Device {
484    /// Get information about the stored history.
485    pub async fn get_history_info(&self) -> Result<HistoryInfo> {
486        // Read total readings count
487        let total_data = self.read_characteristic(TOTAL_READINGS).await?;
488        let total_readings = if total_data.len() >= 2 {
489            u16::from_le_bytes([total_data[0], total_data[1]])
490        } else {
491            return Err(Error::InvalidData(
492                "Invalid total readings data".to_string(),
493            ));
494        };
495
496        // Read interval
497        let interval_data = self.read_characteristic(READ_INTERVAL).await?;
498        let interval_seconds = if interval_data.len() >= 2 {
499            u16::from_le_bytes([interval_data[0], interval_data[1]])
500        } else {
501            return Err(Error::InvalidData("Invalid interval data".to_string()));
502        };
503
504        // Read seconds since update
505        let age_data = self.read_characteristic(SECONDS_SINCE_UPDATE).await?;
506        let seconds_since_update = if age_data.len() >= 2 {
507            u16::from_le_bytes([age_data[0], age_data[1]])
508        } else {
509            0
510        };
511
512        Ok(HistoryInfo {
513            total_readings,
514            interval_seconds,
515            seconds_since_update,
516        })
517    }
518
519    /// Download all historical readings from the device.
520    pub async fn download_history(&self) -> Result<Vec<HistoryRecord>> {
521        self.download_history_with_options(HistoryOptions::default())
522            .await
523    }
524
525    /// Download historical readings with custom options.
526    ///
527    /// # Device Support
528    ///
529    /// - **Aranet4**: Downloads CO₂, temperature, pressure, humidity
530    /// - **Aranet2**: Downloads temperature, humidity
531    /// - **AranetRn+ (Radon)**: Downloads radon, temperature, pressure, humidity
532    /// - **Aranet Radiation**: **Not supported** - returns an error. The device protocol
533    ///   for historical radiation data requires additional documentation. Use
534    ///   [`Device::read_current()`](crate::device::Device::read_current) to get
535    ///   current radiation readings.
536    ///
537    /// # Adaptive Delay
538    ///
539    /// If `options.use_adaptive_delay` is enabled, the read delay will be
540    /// automatically adjusted based on the connection's signal quality.
541    ///
542    /// # Checkpointing
543    ///
544    /// If a checkpoint callback is set, progress will be saved periodically
545    /// to allow resuming interrupted downloads.
546    pub async fn download_history_with_options(
547        &self,
548        options: HistoryOptions,
549    ) -> Result<Vec<HistoryRecord>> {
550        use aranet_types::DeviceType;
551
552        let info = self.get_history_info().await?;
553        // Capture "now" as close as possible to the seconds-since-update read so
554        // every sync maps the same device record to the same timestamp.
555        let anchor = HistoryAnchor::new(&info, OffsetDateTime::now_utc());
556        info!(
557            "Device has {} readings, interval {}s, last update {}s ago",
558            info.total_readings, info.interval_seconds, info.seconds_since_update
559        );
560
561        if info.total_readings == 0 {
562            return Ok(Vec::new());
563        }
564
565        let start_idx = options.start_index.unwrap_or(1);
566        let end_idx = options.end_index.unwrap_or(info.total_readings);
567
568        if start_idx > end_idx {
569            return Err(Error::InvalidConfig(format!(
570                "start_index ({start_idx}) must be <= end_index ({end_idx})"
571            )));
572        }
573        if start_idx == 0 {
574            return Err(Error::InvalidConfig(
575                "start_index must be >= 1 (indices are 1-based)".into(),
576            ));
577        }
578
579        // Get signal quality for adaptive delay if enabled
580        let signal_quality = if options.use_adaptive_delay {
581            match self.signal_quality().await {
582                Some(quality) => {
583                    info!(
584                        "Signal quality: {:?} - using {} ms read delay",
585                        quality,
586                        quality.recommended_read_delay().as_millis()
587                    );
588                    Some(quality)
589                }
590                None => {
591                    debug!("Could not read signal quality, using default delay");
592                    None
593                }
594            }
595        } else {
596            None
597        };
598
599        // Calculate effective read delay
600        let effective_delay = options.effective_read_delay(signal_quality);
601
602        // Dispatch based on device type
603        match self.device_type() {
604            Some(DeviceType::AranetRadiation) => {
605                // Aranet Radiation history download is not supported.
606                // The BLE protocol for historical radiation data differs from other
607                // Aranet devices and is not publicly documented by SAF Tehnika.
608                Err(Error::Unsupported(
609                    "History download is not available for Aranet Radiation devices. \
610                     The radiation history protocol is not documented. \
611                     Use read_current() for current radiation readings."
612                        .to_string(),
613                ))
614            }
615            Some(DeviceType::AranetRadon) => {
616                // For radon devices, download radon instead of CO2, and use Humidity2
617                self.download_radon_history_internal(
618                    &info,
619                    anchor,
620                    start_idx,
621                    end_idx,
622                    &options,
623                    effective_delay,
624                )
625                .await
626            }
627            Some(DeviceType::Aranet2) => {
628                // For Aranet2, download temperature and humidity only
629                self.download_aranet2_history_internal(
630                    &info,
631                    anchor,
632                    start_idx,
633                    end_idx,
634                    &options,
635                    effective_delay,
636                )
637                .await
638            }
639            _ => {
640                // For Aranet4 (and unknown devices), download CO2, temp, pressure, humidity
641                self.download_aranet4_history_internal(
642                    &info,
643                    anchor,
644                    start_idx,
645                    end_idx,
646                    &options,
647                    effective_delay,
648                )
649                .await
650            }
651        }
652    }
653
654    /// Download a u16 parameter with progress reporting and checkpoint updates.
655    ///
656    /// This is the common pattern shared by all parameter downloads except radon (u32).
657    /// Returns the downloaded values.
658    async fn download_u16_param_with_checkpoint(
659        &self,
660        step_info: U16HistoryStep,
661        start_idx: u16,
662        end_idx: u16,
663        effective_delay: Duration,
664        options: &HistoryOptions,
665        checkpoint: &mut Option<HistoryCheckpoint>,
666    ) -> Result<Vec<u16>> {
667        let total_values = (end_idx - start_idx + 1) as usize;
668        let mut progress = HistoryProgress::new(
669            step_info.param,
670            step_info.step,
671            step_info.total_steps,
672            total_values,
673        );
674        options.report_progress(&progress);
675
676        let values = self
677            .download_param_history_with_progress(
678                step_info.param,
679                start_idx,
680                end_idx,
681                effective_delay,
682                |downloaded| {
683                    progress.update(downloaded);
684                    options.report_progress(&progress);
685                },
686            )
687            .await?;
688
689        if let Some(cp) = checkpoint {
690            cp.complete_param(step_info.param, values.clone());
691            if let Some(next) = step_info.next_param {
692                cp.current_param = next;
693                cp.resume_index = start_idx;
694            }
695            options.report_checkpoint(cp);
696        }
697
698        Ok(values)
699    }
700
701    /// Download history for Aranet4 devices (CO2, temp, pressure, humidity).
702    async fn download_aranet4_history_internal(
703        &self,
704        info: &HistoryInfo,
705        anchor: HistoryAnchor,
706        start_idx: u16,
707        end_idx: u16,
708        options: &HistoryOptions,
709        effective_delay: Duration,
710    ) -> Result<Vec<HistoryRecord>> {
711        if start_idx > end_idx {
712            return Ok(Vec::new());
713        }
714
715        let device_id = self.address().to_string();
716        let mut checkpoint = if options.checkpoint_callback.is_some() {
717            Some(HistoryCheckpoint::new(
718                &device_id,
719                info.total_readings,
720                HistoryParam::Co2,
721            ))
722        } else {
723            None
724        };
725
726        let co2_values = self
727            .download_u16_param_with_checkpoint(
728                U16HistoryStep {
729                    param: HistoryParam::Co2,
730                    step: 1,
731                    total_steps: 4,
732                    next_param: Some(HistoryParamCheckpoint::Temperature),
733                },
734                start_idx,
735                end_idx,
736                effective_delay,
737                options,
738                &mut checkpoint,
739            )
740            .await?;
741
742        let temp_values = self
743            .download_u16_param_with_checkpoint(
744                U16HistoryStep {
745                    param: HistoryParam::Temperature,
746                    step: 2,
747                    total_steps: 4,
748                    next_param: Some(HistoryParamCheckpoint::Pressure),
749                },
750                start_idx,
751                end_idx,
752                effective_delay,
753                options,
754                &mut checkpoint,
755            )
756            .await?;
757
758        let pressure_values = self
759            .download_u16_param_with_checkpoint(
760                U16HistoryStep {
761                    param: HistoryParam::Pressure,
762                    step: 3,
763                    total_steps: 4,
764                    next_param: Some(HistoryParamCheckpoint::Humidity),
765                },
766                start_idx,
767                end_idx,
768                effective_delay,
769                options,
770                &mut checkpoint,
771            )
772            .await?;
773
774        let humidity_values = self
775            .download_u16_param_with_checkpoint(
776                U16HistoryStep {
777                    param: HistoryParam::Humidity,
778                    step: 4,
779                    total_steps: 4,
780                    next_param: None,
781                },
782                start_idx,
783                end_idx,
784                effective_delay,
785                options,
786                &mut checkpoint,
787            )
788            .await?;
789
790        let records = build_history_records(
791            &anchor,
792            start_idx,
793            &co2_values,
794            &temp_values,
795            &pressure_values,
796            &humidity_values,
797            &[],
798        );
799
800        info!("Downloaded {} history records", records.len());
801        Ok(records)
802    }
803
804    /// Download history for Aranet2 devices (temperature, humidity only).
805    async fn download_aranet2_history_internal(
806        &self,
807        info: &HistoryInfo,
808        anchor: HistoryAnchor,
809        start_idx: u16,
810        end_idx: u16,
811        options: &HistoryOptions,
812        effective_delay: Duration,
813    ) -> Result<Vec<HistoryRecord>> {
814        if start_idx > end_idx {
815            return Ok(Vec::new());
816        }
817
818        let device_id = self.address().to_string();
819        let mut checkpoint = if options.checkpoint_callback.is_some() {
820            Some(HistoryCheckpoint::new(
821                &device_id,
822                info.total_readings,
823                HistoryParam::Temperature,
824            ))
825        } else {
826            None
827        };
828
829        let temp_values = self
830            .download_u16_param_with_checkpoint(
831                U16HistoryStep {
832                    param: HistoryParam::Temperature,
833                    step: 1,
834                    total_steps: 2,
835                    next_param: Some(HistoryParamCheckpoint::Humidity2),
836                },
837                start_idx,
838                end_idx,
839                effective_delay,
840                options,
841                &mut checkpoint,
842            )
843            .await?;
844
845        let humidity_values = self
846            .download_u16_param_with_checkpoint(
847                U16HistoryStep {
848                    param: HistoryParam::Humidity2,
849                    step: 2,
850                    total_steps: 2,
851                    next_param: None,
852                },
853                start_idx,
854                end_idx,
855                effective_delay,
856                options,
857                &mut checkpoint,
858            )
859            .await?;
860
861        // Build records with no CO2, no pressure, no radon
862        let records = build_history_records(
863            &anchor,
864            start_idx,
865            &[],
866            &temp_values,
867            &[],
868            &humidity_values,
869            &[],
870        );
871
872        info!("Downloaded {} Aranet2 history records", records.len());
873        Ok(records)
874    }
875
876    /// Download history for AranetRn+ devices (radon, temp, pressure, humidity).
877    async fn download_radon_history_internal(
878        &self,
879        info: &HistoryInfo,
880        anchor: HistoryAnchor,
881        start_idx: u16,
882        end_idx: u16,
883        options: &HistoryOptions,
884        effective_delay: Duration,
885    ) -> Result<Vec<HistoryRecord>> {
886        if start_idx > end_idx {
887            return Ok(Vec::new());
888        }
889        let total_values = (end_idx - start_idx + 1) as usize;
890
891        let device_id = self.address().to_string();
892        let mut checkpoint = if options.checkpoint_callback.is_some() {
893            Some(HistoryCheckpoint::new(
894                &device_id,
895                info.total_readings,
896                HistoryParam::Radon,
897            ))
898        } else {
899            None
900        };
901
902        // Download radon values (4 bytes each, uses u32 variant)
903        let mut progress = HistoryProgress::new(HistoryParam::Radon, 1, 4, total_values);
904        options.report_progress(&progress);
905
906        let radon_values = self
907            .download_param_history_u32_with_progress(
908                HistoryParam::Radon,
909                start_idx,
910                end_idx,
911                effective_delay,
912                |downloaded| {
913                    progress.update(downloaded);
914                    options.report_progress(&progress);
915                },
916            )
917            .await?;
918
919        if let Some(ref mut cp) = checkpoint {
920            cp.complete_radon_param(radon_values.clone());
921            cp.current_param = HistoryParamCheckpoint::Temperature;
922            cp.resume_index = start_idx;
923            options.report_checkpoint(cp);
924        }
925
926        let temp_values = self
927            .download_u16_param_with_checkpoint(
928                U16HistoryStep {
929                    param: HistoryParam::Temperature,
930                    step: 2,
931                    total_steps: 4,
932                    next_param: Some(HistoryParamCheckpoint::Pressure),
933                },
934                start_idx,
935                end_idx,
936                effective_delay,
937                options,
938                &mut checkpoint,
939            )
940            .await?;
941
942        let pressure_values = self
943            .download_u16_param_with_checkpoint(
944                U16HistoryStep {
945                    param: HistoryParam::Pressure,
946                    step: 3,
947                    total_steps: 4,
948                    next_param: Some(HistoryParamCheckpoint::Humidity2),
949                },
950                start_idx,
951                end_idx,
952                effective_delay,
953                options,
954                &mut checkpoint,
955            )
956            .await?;
957
958        let humidity_values = self
959            .download_u16_param_with_checkpoint(
960                U16HistoryStep {
961                    param: HistoryParam::Humidity2,
962                    step: 4,
963                    total_steps: 4,
964                    next_param: None,
965                },
966                start_idx,
967                end_idx,
968                effective_delay,
969                options,
970                &mut checkpoint,
971            )
972            .await?;
973
974        let records = build_history_records(
975            &anchor,
976            start_idx,
977            &[],
978            &temp_values,
979            &pressure_values,
980            &humidity_values,
981            &radon_values,
982        );
983
984        info!("Downloaded {} radon history records", records.len());
985        Ok(records)
986    }
987
988    /// Download a single parameter's history using V2 protocol with progress callback.
989    ///
990    /// This is a generic implementation that handles different value sizes:
991    /// - 1 byte: humidity
992    /// - 2 bytes: CO2, temperature, pressure, humidity2
993    /// - 4 bytes: radon
994    #[allow(clippy::too_many_arguments)]
995    async fn download_param_history_generic_with_progress<T, F>(
996        &self,
997        param: HistoryParam,
998        start_idx: u16,
999        end_idx: u16,
1000        read_delay: Duration,
1001        value_parser: impl Fn(&[u8], usize) -> Option<T>,
1002        value_size: usize,
1003        mut on_progress: F,
1004    ) -> Result<Vec<T>>
1005    where
1006        T: Default + Clone,
1007        F: FnMut(usize),
1008    {
1009        debug!(
1010            "Downloading {:?} history from {} to {} (value_size={})",
1011            param, start_idx, end_idx, value_size
1012        );
1013
1014        let mut values: BTreeMap<u16, T> = BTreeMap::new();
1015        let mut current_idx = start_idx;
1016        let mut consecutive_wrong_param = 0u32;
1017        const MAX_WRONG_PARAM_RETRIES: u32 = 5;
1018        let mut consecutive_stalls = 0u32;
1019        const MAX_STALLED_PACKETS: u32 = 5;
1020
1021        while current_idx <= end_idx {
1022            // Send V2 history request using command constant
1023            let cmd = [
1024                HISTORY_V2_REQUEST,
1025                param as u8,
1026                (current_idx & 0xFF) as u8,
1027                ((current_idx >> 8) & 0xFF) as u8,
1028            ];
1029
1030            self.write_characteristic(COMMAND, &cmd).await?;
1031            sleep(read_delay).await;
1032
1033            // Read response
1034            let response = self.read_characteristic(HISTORY_V2).await?;
1035
1036            // V2 response format (10-byte header):
1037            // Byte 0: param (1 byte)
1038            // Bytes 1-2: interval (2 bytes, little-endian)
1039            // Bytes 3-4: total_readings (2 bytes, little-endian)
1040            // Bytes 5-6: ago (2 bytes, little-endian)
1041            // Bytes 7-8: start index (2 bytes, little-endian)
1042            // Byte 9: count (1 byte)
1043            // Bytes 10+: data values
1044            if response.len() < 10 {
1045                warn!(
1046                    "Invalid history response: too short ({} bytes)",
1047                    response.len()
1048                );
1049                break;
1050            }
1051
1052            let resp_param = response[0];
1053            if resp_param != param as u8 {
1054                consecutive_wrong_param += 1;
1055                warn!(
1056                    "Unexpected parameter in response: {} (retry {}/{})",
1057                    resp_param, consecutive_wrong_param, MAX_WRONG_PARAM_RETRIES
1058                );
1059                if consecutive_wrong_param >= MAX_WRONG_PARAM_RETRIES {
1060                    warn!("Too many wrong parameter responses, aborting download");
1061                    break;
1062                }
1063                // Wait and retry - device may not have processed command yet
1064                sleep(read_delay).await;
1065                continue;
1066            }
1067            consecutive_wrong_param = 0;
1068
1069            // Parse header
1070            let resp_start = u16::from_le_bytes([response[7], response[8]]);
1071            let resp_count = response[9] as usize;
1072
1073            debug!(
1074                "History response: param={}, start={}, count={}",
1075                resp_param, resp_start, resp_count
1076            );
1077
1078            // Check if we've reached the end (count == 0)
1079            if resp_count == 0 {
1080                debug!("Reached end of history (count=0)");
1081                break;
1082            }
1083
1084            // Parse data values and decide what to request next
1085            let data = &response[10..];
1086            match apply_v2_packet(
1087                &mut values,
1088                current_idx,
1089                end_idx,
1090                resp_start,
1091                resp_count,
1092                data,
1093                value_size,
1094                &value_parser,
1095            ) {
1096                PacketProgress::Done => {
1097                    on_progress(values.len());
1098                    debug!("Reached end of requested range");
1099                    break;
1100                }
1101                PacketProgress::Continue(next_idx) => {
1102                    consecutive_stalls = 0;
1103                    current_idx = next_idx;
1104                    debug!(
1105                        "Downloaded {} values, next index: {}",
1106                        values.len(),
1107                        current_idx
1108                    );
1109                    on_progress(values.len());
1110                }
1111                PacketProgress::Stalled => {
1112                    consecutive_stalls += 1;
1113                    warn!(
1114                        "History packet for {:?} made no progress at index {} (retry {}/{})",
1115                        param, current_idx, consecutive_stalls, MAX_STALLED_PACKETS
1116                    );
1117                    if consecutive_stalls >= MAX_STALLED_PACKETS {
1118                        return Err(Error::InvalidData(format!(
1119                            "History download for {param:?} stalled at index {current_idx}"
1120                        )));
1121                    }
1122                    sleep(read_delay).await;
1123                }
1124            }
1125        }
1126
1127        // Convert to ordered vector (BTreeMap already maintains order)
1128        Ok(values.into_values().collect())
1129    }
1130
1131    /// Download a single parameter's history using V2 protocol (u16 values) with progress.
1132    async fn download_param_history_with_progress<F>(
1133        &self,
1134        param: HistoryParam,
1135        start_idx: u16,
1136        end_idx: u16,
1137        read_delay: Duration,
1138        on_progress: F,
1139    ) -> Result<Vec<u16>>
1140    where
1141        F: FnMut(usize),
1142    {
1143        let value_size = if param == HistoryParam::Humidity {
1144            1
1145        } else {
1146            2
1147        };
1148
1149        self.download_param_history_generic_with_progress(
1150            param,
1151            start_idx,
1152            end_idx,
1153            read_delay,
1154            |data, i| {
1155                if param == HistoryParam::Humidity {
1156                    data.get(i).map(|&b| b as u16)
1157                } else {
1158                    let offset = i * 2;
1159                    if offset + 1 < data.len() {
1160                        Some(u16::from_le_bytes([data[offset], data[offset + 1]]))
1161                    } else {
1162                        None
1163                    }
1164                }
1165            },
1166            value_size,
1167            on_progress,
1168        )
1169        .await
1170    }
1171
1172    /// Download a single parameter's history using V2 protocol (u32 values) with progress.
1173    async fn download_param_history_u32_with_progress<F>(
1174        &self,
1175        param: HistoryParam,
1176        start_idx: u16,
1177        end_idx: u16,
1178        read_delay: Duration,
1179        on_progress: F,
1180    ) -> Result<Vec<u32>>
1181    where
1182        F: FnMut(usize),
1183    {
1184        self.download_param_history_generic_with_progress(
1185            param,
1186            start_idx,
1187            end_idx,
1188            read_delay,
1189            |data, i| {
1190                let offset = i * 4;
1191                if offset + 3 < data.len() {
1192                    Some(u32::from_le_bytes([
1193                        data[offset],
1194                        data[offset + 1],
1195                        data[offset + 2],
1196                        data[offset + 3],
1197                    ]))
1198                } else {
1199                    None
1200                }
1201            },
1202            4,
1203            on_progress,
1204        )
1205        .await
1206    }
1207
1208    /// Download history using V1 protocol (notification-based).
1209    ///
1210    /// This is used for older devices that don't support the V2 read-based protocol.
1211    /// V1 uses notifications on the HISTORY_V1 characteristic.
1212    pub async fn download_history_v1(&self) -> Result<Vec<HistoryRecord>> {
1213        use crate::uuid::HISTORY_V1;
1214        use tokio::sync::mpsc;
1215
1216        let info = self.get_history_info().await?;
1217        info!(
1218            "V1 download: {} readings, interval {}s",
1219            info.total_readings, info.interval_seconds
1220        );
1221
1222        if info.total_readings == 0 {
1223            return Ok(Vec::new());
1224        }
1225
1226        // Subscribe to notifications
1227        let (tx, mut rx) = mpsc::channel::<Vec<u8>>(256);
1228
1229        // Set up notification handler
1230        self.subscribe_to_notifications(HISTORY_V1, move |data| {
1231            if let Err(e) = tx.try_send(data.to_vec()) {
1232                warn!(
1233                    "V1 history notification channel full or closed, data may be lost: {}",
1234                    e
1235                );
1236            }
1237        })
1238        .await?;
1239
1240        // Request history for each parameter
1241        let mut co2_values = Vec::new();
1242        let mut temp_values = Vec::new();
1243        let mut pressure_values = Vec::new();
1244        let mut humidity_values = Vec::new();
1245
1246        for param in [
1247            HistoryParam::Co2,
1248            HistoryParam::Temperature,
1249            HistoryParam::Pressure,
1250            HistoryParam::Humidity,
1251        ] {
1252            // Send V1 history request using command constant
1253            let cmd = [
1254                HISTORY_V1_REQUEST,
1255                param as u8,
1256                0x01,
1257                0x00,
1258                (info.total_readings & 0xFF) as u8,
1259                ((info.total_readings >> 8) & 0xFF) as u8,
1260            ];
1261
1262            if let Err(e) = self.write_characteristic(COMMAND, &cmd).await {
1263                // The unsubscribe after the loop is skipped on this path; stop
1264                // the notification task here. The write error is the one to
1265                // report.
1266                let _ = self.unsubscribe_from_notifications(HISTORY_V1).await;
1267                return Err(e);
1268            }
1269
1270            // Collect notifications until we have all values
1271            let mut values = Vec::new();
1272            let expected = info.total_readings as usize;
1273
1274            let mut consecutive_timeouts = 0;
1275            const MAX_CONSECUTIVE_TIMEOUTS: u32 = 3;
1276
1277            while values.len() < expected {
1278                match tokio::time::timeout(Duration::from_secs(5), rx.recv()).await {
1279                    Ok(Some(data)) => {
1280                        consecutive_timeouts = 0; // Reset on successful receive
1281                        // Parse notification data
1282                        if data.len() >= 3 {
1283                            let resp_param = data[0];
1284                            if resp_param == param as u8 {
1285                                let mut buf = &data[3..];
1286                                while buf.len() >= 2 && values.len() < expected {
1287                                    values.push(buf.get_u16_le());
1288                                }
1289                            }
1290                        }
1291                    }
1292                    Ok(None) => {
1293                        warn!(
1294                            "V1 history channel closed for {:?}: got {}/{} values",
1295                            param,
1296                            values.len(),
1297                            expected
1298                        );
1299                        break;
1300                    }
1301                    Err(_) => {
1302                        consecutive_timeouts += 1;
1303                        warn!(
1304                            "Timeout waiting for V1 history notification ({}/{}), {:?}: {}/{} values",
1305                            consecutive_timeouts,
1306                            MAX_CONSECUTIVE_TIMEOUTS,
1307                            param,
1308                            values.len(),
1309                            expected
1310                        );
1311                        if consecutive_timeouts >= MAX_CONSECUTIVE_TIMEOUTS {
1312                            warn!(
1313                                "Too many consecutive timeouts for {:?}, proceeding with partial data",
1314                                param
1315                            );
1316                            break;
1317                        }
1318                    }
1319                }
1320            }
1321
1322            // Log if we got incomplete data
1323            if values.len() < expected {
1324                warn!(
1325                    "V1 history download incomplete for {:?}: got {}/{} values ({:.1}%)",
1326                    param,
1327                    values.len(),
1328                    expected,
1329                    (values.len() as f64 / expected as f64) * 100.0
1330                );
1331            }
1332
1333            match param {
1334                HistoryParam::Co2 => co2_values = values,
1335                HistoryParam::Temperature => temp_values = values,
1336                HistoryParam::Pressure => pressure_values = values,
1337                HistoryParam::Humidity => humidity_values = values,
1338                // V1 protocol doesn't support radon or humidity2
1339                HistoryParam::Humidity2 | HistoryParam::Radon => {}
1340            }
1341        }
1342
1343        // Unsubscribe from notifications
1344        self.unsubscribe_from_notifications(HISTORY_V1).await?;
1345
1346        // Build history records
1347        let now = OffsetDateTime::now_utc();
1348        let latest_reading_time = now - time::Duration::seconds(info.seconds_since_update as i64);
1349
1350        let mut records = Vec::new();
1351        let count = co2_values.len();
1352
1353        // Warn if parameter arrays have mismatched lengths (partial download)
1354        if temp_values.len() != count
1355            || pressure_values.len() != count
1356            || humidity_values.len() != count
1357        {
1358            warn!(
1359                "V1 history arrays have mismatched lengths: co2={}, temp={}, pressure={}, humidity={} — \
1360                 records with missing values will use defaults",
1361                count,
1362                temp_values.len(),
1363                pressure_values.len(),
1364                humidity_values.len()
1365            );
1366        }
1367
1368        for i in 0..count {
1369            let readings_ago = (count - 1 - i) as i64;
1370            let timestamp = latest_reading_time
1371                - time::Duration::seconds(readings_ago * info.interval_seconds as i64);
1372
1373            let record = HistoryRecord {
1374                timestamp,
1375                co2: co2_values.get(i).copied().unwrap_or(0),
1376                temperature: raw_to_temperature(temp_values.get(i).copied().unwrap_or(0)),
1377                pressure: raw_to_pressure(pressure_values.get(i).copied().unwrap_or(0)),
1378                humidity: humidity_values.get(i).copied().unwrap_or(0) as u8,
1379                radon: None,
1380                radiation_rate: None,
1381                radiation_total: None,
1382            };
1383            records.push(record);
1384        }
1385
1386        info!("V1 download complete: {} records", records.len());
1387        Ok(records)
1388    }
1389}
1390
1391/// Outcome of applying one V2 history response packet.
1392#[derive(Debug, PartialEq, Eq)]
1393enum PacketProgress {
1394    /// More records remain; request the next packet starting at this index.
1395    Continue(u16),
1396    /// The packet reached `end_idx`; there is nothing more to request.
1397    Done,
1398    /// The packet did not move the download forward (stale repeat, or a
1399    /// count with no payload). The caller should retry, then give up.
1400    Stalled,
1401}
1402
1403/// Store the values from one V2 history packet and decide what to request next.
1404///
1405/// `resp_start` and `resp_count` come from the packet header; `data` is the
1406/// payload after the 10-byte header. `end_idx` is inclusive.
1407#[allow(clippy::too_many_arguments)]
1408fn apply_v2_packet<T>(
1409    values: &mut BTreeMap<u16, T>,
1410    current_idx: u16,
1411    end_idx: u16,
1412    resp_start: u16,
1413    resp_count: usize,
1414    data: &[u8],
1415    value_size: usize,
1416    value_parser: &impl Fn(&[u8], usize) -> Option<T>,
1417) -> PacketProgress {
1418    let num_values = (data.len() / value_size).min(resp_count);
1419    if num_values == 0 {
1420        return PacketProgress::Stalled;
1421    }
1422
1423    for i in 0..num_values {
1424        let Some(idx) = resp_start.checked_add(i as u16) else {
1425            break;
1426        };
1427        if idx > end_idx {
1428            break;
1429        }
1430        // A stale packet can start below the index we asked for. Keep only the
1431        // requested range: callers stamp element i as device index start_idx + i.
1432        if idx < current_idx {
1433            continue;
1434        }
1435        if let Some(value) = value_parser(data, i) {
1436            values.insert(idx, value);
1437        }
1438    }
1439
1440    // Index of the last value carried by this packet (num_values >= 1).
1441    let last_idx = u32::from(resp_start) + num_values as u32 - 1;
1442    if last_idx >= u32::from(end_idx) {
1443        return PacketProgress::Done;
1444    }
1445    // last_idx < end_idx <= u16::MAX, so this fits in u16.
1446    let next_idx = (last_idx + 1) as u16;
1447    if next_idx <= current_idx {
1448        return PacketProgress::Stalled;
1449    }
1450    PacketProgress::Continue(next_idx)
1451}
1452
1453/// Build history records from downloaded parameter arrays.
1454///
1455/// Element `i` of each array is the record at 1-based device index
1456/// `start_idx + i`; its timestamp comes from `anchor` for that index, so a
1457/// partial download (e.g. `end_index < total_readings`) keeps the device's
1458/// real record times instead of being stamped as the newest records.
1459///
1460/// For Aranet4: pass co2_values and empty radon_values.
1461/// For AranetRn+: pass empty co2_values and radon_values.
1462/// Humidity is converted differently based on whether radon_values is populated
1463/// (radon devices use Humidity2 encoding: tenths of a percent).
1464fn build_history_records(
1465    anchor: &HistoryAnchor,
1466    start_idx: u16,
1467    co2_values: &[u16],
1468    temp_values: &[u16],
1469    pressure_values: &[u16],
1470    humidity_values: &[u16],
1471    radon_values: &[u32],
1472) -> Vec<HistoryRecord> {
1473    let is_radon = !radon_values.is_empty();
1474    let is_aranet2 = co2_values.is_empty() && radon_values.is_empty();
1475    let count = if is_radon {
1476        radon_values.len()
1477    } else if is_aranet2 {
1478        temp_values.len()
1479    } else {
1480        co2_values.len()
1481    };
1482
1483    // Warn if parameter arrays have mismatched lengths (partial download)
1484    let expected = count;
1485    if temp_values.len() != expected
1486        || pressure_values.len() != expected
1487        || humidity_values.len() != expected
1488    {
1489        warn!(
1490            "History arrays have mismatched lengths: primary={expected}, temp={}, pressure={}, humidity={} — \
1491             records with missing values will use defaults",
1492            temp_values.len(),
1493            pressure_values.len(),
1494            humidity_values.len()
1495        );
1496    }
1497
1498    (0..count)
1499        .map(|i| {
1500            let timestamp = anchor.timestamp_for(u32::from(start_idx) + i as u32);
1501
1502            let humidity = if is_radon || is_aranet2 {
1503                // Humidity2 is stored as tenths of a percent
1504                let raw = humidity_values.get(i).copied().unwrap_or(0);
1505                (raw / 10).min(100) as u8
1506            } else {
1507                humidity_values.get(i).copied().unwrap_or(0) as u8
1508            };
1509
1510            HistoryRecord {
1511                timestamp,
1512                co2: if is_radon {
1513                    0
1514                } else {
1515                    co2_values.get(i).copied().unwrap_or(0)
1516                },
1517                temperature: raw_to_temperature(temp_values.get(i).copied().unwrap_or(0)),
1518                pressure: raw_to_pressure(pressure_values.get(i).copied().unwrap_or(0)),
1519                humidity,
1520                radon: if is_radon {
1521                    Some(radon_values.get(i).copied().unwrap_or(0))
1522                } else {
1523                    None
1524                },
1525                radiation_rate: None,
1526                radiation_total: None,
1527            }
1528        })
1529        .collect()
1530}
1531
1532/// Convert raw temperature value to Celsius.
1533pub fn raw_to_temperature(raw: u16) -> f32 {
1534    raw as f32 / 20.0
1535}
1536
1537/// Convert raw pressure value to hPa.
1538pub fn raw_to_pressure(raw: u16) -> f32 {
1539    raw as f32 / 10.0
1540}
1541
1542// NOTE: The HistoryValueConverter trait was removed as it was dead code.
1543// Use the standalone functions raw_to_temperature, raw_to_pressure, etc. directly.
1544
1545#[cfg(test)]
1546mod tests {
1547    use super::*;
1548
1549    // --- raw_to_temperature tests ---
1550
1551    #[test]
1552    fn test_raw_to_temperature_typical_values() {
1553        // 22.5°C = 450 raw (450/20 = 22.5)
1554        assert!((raw_to_temperature(450) - 22.5).abs() < 0.001);
1555
1556        // 20.0°C = 400 raw
1557        assert!((raw_to_temperature(400) - 20.0).abs() < 0.001);
1558
1559        // 25.0°C = 500 raw
1560        assert!((raw_to_temperature(500) - 25.0).abs() < 0.001);
1561    }
1562
1563    #[test]
1564    fn test_raw_to_temperature_edge_cases() {
1565        // 0°C = 0 raw
1566        assert!((raw_to_temperature(0) - 0.0).abs() < 0.001);
1567
1568        // Very cold: -10°C would be negative, but raw is u16 so minimum is 0
1569        // Raw values represent actual temperature * 20
1570
1571        // Very hot: 50°C = 1000 raw
1572        assert!((raw_to_temperature(1000) - 50.0).abs() < 0.001);
1573
1574        // Maximum u16 would be 65535/20 = 3276.75°C (unrealistic but tests overflow handling)
1575        assert!((raw_to_temperature(u16::MAX) - 3276.75).abs() < 0.01);
1576    }
1577
1578    #[test]
1579    fn test_raw_to_temperature_precision() {
1580        // Test fractional values
1581        // 22.55°C = 451 raw
1582        assert!((raw_to_temperature(451) - 22.55).abs() < 0.001);
1583
1584        // 22.05°C = 441 raw
1585        assert!((raw_to_temperature(441) - 22.05).abs() < 0.001);
1586    }
1587
1588    // --- raw_to_pressure tests ---
1589
1590    #[test]
1591    fn test_raw_to_pressure_typical_values() {
1592        // 1013.2 hPa = 10132 raw
1593        assert!((raw_to_pressure(10132) - 1013.2).abs() < 0.01);
1594
1595        // 1000.0 hPa = 10000 raw
1596        assert!((raw_to_pressure(10000) - 1000.0).abs() < 0.01);
1597
1598        // 1050.0 hPa = 10500 raw
1599        assert!((raw_to_pressure(10500) - 1050.0).abs() < 0.01);
1600    }
1601
1602    #[test]
1603    fn test_raw_to_pressure_edge_cases() {
1604        // 0 hPa = 0 raw
1605        assert!((raw_to_pressure(0) - 0.0).abs() < 0.01);
1606
1607        // Low pressure: 950 hPa = 9500 raw
1608        assert!((raw_to_pressure(9500) - 950.0).abs() < 0.01);
1609
1610        // High pressure: 1100 hPa = 11000 raw
1611        assert!((raw_to_pressure(11000) - 1100.0).abs() < 0.01);
1612
1613        // Maximum u16 would be 65535/10 = 6553.5 hPa (unrealistic but tests bounds)
1614        assert!((raw_to_pressure(u16::MAX) - 6553.5).abs() < 0.1);
1615    }
1616
1617    // --- HistoryParam tests ---
1618
1619    #[test]
1620    fn test_history_param_values() {
1621        assert_eq!(HistoryParam::Temperature as u8, 1);
1622        assert_eq!(HistoryParam::Humidity as u8, 2);
1623        assert_eq!(HistoryParam::Pressure as u8, 3);
1624        assert_eq!(HistoryParam::Co2 as u8, 4);
1625    }
1626
1627    #[test]
1628    fn test_history_param_debug() {
1629        assert_eq!(format!("{:?}", HistoryParam::Temperature), "Temperature");
1630        assert_eq!(format!("{:?}", HistoryParam::Co2), "Co2");
1631    }
1632
1633    // --- HistoryOptions tests ---
1634
1635    #[test]
1636    fn test_history_options_default() {
1637        let options = HistoryOptions::default();
1638
1639        assert!(options.start_index.is_none());
1640        assert!(options.end_index.is_none());
1641        assert_eq!(options.read_delay, Duration::from_millis(50));
1642    }
1643
1644    #[test]
1645    fn test_history_options_custom() {
1646        let options = HistoryOptions::new()
1647            .start_index(10)
1648            .end_index(100)
1649            .read_delay(Duration::from_millis(100));
1650
1651        assert_eq!(options.start_index, Some(10));
1652        assert_eq!(options.end_index, Some(100));
1653        assert_eq!(options.read_delay, Duration::from_millis(100));
1654    }
1655
1656    #[test]
1657    fn test_history_options_with_progress() {
1658        use std::sync::Arc;
1659        use std::sync::atomic::{AtomicUsize, Ordering};
1660
1661        let call_count = Arc::new(AtomicUsize::new(0));
1662        let call_count_clone = Arc::clone(&call_count);
1663
1664        let options = HistoryOptions::new().with_progress(move |_progress| {
1665            call_count_clone.fetch_add(1, Ordering::SeqCst);
1666        });
1667
1668        assert!(options.progress_callback.is_some());
1669
1670        // Test that the callback can be invoked
1671        let progress = HistoryProgress::new(HistoryParam::Co2, 1, 4, 100);
1672        options.report_progress(&progress);
1673        assert_eq!(call_count.load(Ordering::SeqCst), 1);
1674    }
1675
1676    // --- HistoryInfo tests ---
1677
1678    #[test]
1679    fn test_history_info_creation() {
1680        let info = HistoryInfo {
1681            total_readings: 1000,
1682            interval_seconds: 300,
1683            seconds_since_update: 120,
1684        };
1685
1686        assert_eq!(info.total_readings, 1000);
1687        assert_eq!(info.interval_seconds, 300);
1688        assert_eq!(info.seconds_since_update, 120);
1689    }
1690
1691    #[test]
1692    fn test_history_info_debug() {
1693        let info = HistoryInfo {
1694            total_readings: 500,
1695            interval_seconds: 60,
1696            seconds_since_update: 30,
1697        };
1698
1699        let debug_str = format!("{:?}", info);
1700        assert!(debug_str.contains("total_readings"));
1701        assert!(debug_str.contains("500"));
1702    }
1703
1704    // --- apply_v2_packet tests ---
1705
1706    fn parse_u16(data: &[u8], i: usize) -> Option<u16> {
1707        data.get(i * 2..i * 2 + 2)
1708            .map(|b| u16::from_le_bytes([b[0], b[1]]))
1709    }
1710
1711    fn u16_payload(values: &[u16]) -> Vec<u8> {
1712        values.iter().flat_map(|v| v.to_le_bytes()).collect()
1713    }
1714
1715    #[test]
1716    fn test_apply_v2_packet_does_not_stop_one_record_early() {
1717        let mut values = BTreeMap::new();
1718        let payload = u16_payload(&(1..=99).collect::<Vec<u16>>());
1719        let progress = apply_v2_packet(&mut values, 1, 100, 1, 99, &payload, 2, &parse_u16);
1720        assert_eq!(progress, PacketProgress::Continue(100));
1721        assert_eq!(values.len(), 99);
1722    }
1723
1724    #[test]
1725    fn test_apply_v2_packet_done_after_last_index() {
1726        let mut values = BTreeMap::new();
1727        let payload = u16_payload(&[42]);
1728        let progress = apply_v2_packet(&mut values, 100, 100, 100, 1, &payload, 2, &parse_u16);
1729        assert_eq!(progress, PacketProgress::Done);
1730        assert_eq!(values.get(&100), Some(&42));
1731    }
1732
1733    #[test]
1734    fn test_apply_v2_packet_ignores_values_past_end() {
1735        let mut values = BTreeMap::new();
1736        let payload = u16_payload(&[1, 2, 3, 4]);
1737        let progress = apply_v2_packet(&mut values, 9, 10, 9, 4, &payload, 2, &parse_u16);
1738        assert_eq!(progress, PacketProgress::Done);
1739        assert_eq!(values.keys().copied().collect::<Vec<_>>(), vec![9, 10]);
1740    }
1741
1742    #[test]
1743    fn test_apply_v2_packet_count_without_payload_is_stalled() {
1744        let mut values: BTreeMap<u16, u16> = BTreeMap::new();
1745        let progress = apply_v2_packet(&mut values, 1, 100, 1, 5, &[], 2, &parse_u16);
1746        assert_eq!(progress, PacketProgress::Stalled);
1747    }
1748
1749    #[test]
1750    fn test_apply_v2_packet_repeated_old_packet_is_stalled() {
1751        let mut values = BTreeMap::new();
1752        let payload = u16_payload(&[1, 2, 3]);
1753        // We asked for index 50; the device repeated the packet for 1..=3.
1754        let progress = apply_v2_packet(&mut values, 50, 100, 1, 3, &payload, 2, &parse_u16);
1755        assert_eq!(progress, PacketProgress::Stalled);
1756        // Values from below the requested range must not be kept: element i of
1757        // the downloaded array is stamped as device index start_idx + i.
1758        assert!(values.is_empty());
1759    }
1760
1761    #[test]
1762    fn test_apply_v2_packet_same_packet_twice_is_stalled() {
1763        let mut values = BTreeMap::new();
1764        let payload = u16_payload(&[1, 2, 3]);
1765        let progress = apply_v2_packet(&mut values, 1, 100, 1, 3, &payload, 2, &parse_u16);
1766        assert_eq!(progress, PacketProgress::Continue(4));
1767        // Asked for index 4, the device sent the same packet (1..=3) again.
1768        let progress = apply_v2_packet(&mut values, 4, 100, 1, 3, &payload, 2, &parse_u16);
1769        assert_eq!(progress, PacketProgress::Stalled);
1770        assert_eq!(values.keys().copied().collect::<Vec<_>>(), vec![1, 2, 3]);
1771    }
1772
1773    #[test]
1774    fn test_apply_v2_packet_drops_values_below_current_index() {
1775        let mut values = BTreeMap::new();
1776        let payload = u16_payload(&(1..=60).collect::<Vec<u16>>());
1777        // Asked for index 50; a stale packet starting at 1 overlaps the range.
1778        let progress = apply_v2_packet(&mut values, 50, 100, 1, 60, &payload, 2, &parse_u16);
1779        assert_eq!(progress, PacketProgress::Continue(61));
1780        assert_eq!(
1781            values.keys().copied().collect::<Vec<_>>(),
1782            (50..=60).collect::<Vec<u16>>()
1783        );
1784        assert_eq!(values.get(&50), Some(&50));
1785    }
1786
1787    #[test]
1788    fn test_apply_v2_packet_stale_packet_past_end_keeps_only_requested_range() {
1789        let mut values = BTreeMap::new();
1790        let payload = u16_payload(&(1..=117).collect::<Vec<u16>>());
1791        let progress = apply_v2_packet(&mut values, 50, 55, 1, 117, &payload, 2, &parse_u16);
1792        assert_eq!(progress, PacketProgress::Done);
1793        assert_eq!(
1794            values.keys().copied().collect::<Vec<_>>(),
1795            (50..=55).collect::<Vec<u16>>()
1796        );
1797    }
1798
1799    #[test]
1800    fn test_apply_v2_packet_near_u16_max_does_not_overflow() {
1801        let mut values = BTreeMap::new();
1802        let payload = u16_payload(&[1, 2, 3, 4, 5]);
1803        let progress = apply_v2_packet(
1804            &mut values,
1805            65534,
1806            u16::MAX,
1807            65534,
1808            5,
1809            &payload,
1810            2,
1811            &parse_u16,
1812        );
1813        assert_eq!(progress, PacketProgress::Done);
1814        assert_eq!(values.len(), 2);
1815    }
1816
1817    // --- HistoryAnchor / build_history_records timestamp tests ---
1818
1819    fn anchor_fixture(total: u16, interval: u16, ago: u16) -> (HistoryAnchor, OffsetDateTime) {
1820        let read_at = OffsetDateTime::from_unix_timestamp(1_700_000_000).unwrap();
1821        let info = HistoryInfo {
1822            total_readings: total,
1823            interval_seconds: interval,
1824            seconds_since_update: ago,
1825        };
1826        (HistoryAnchor::new(&info, read_at), read_at)
1827    }
1828
1829    #[test]
1830    fn test_history_anchor_newest_record_time() {
1831        let (anchor, read_at) = anchor_fixture(100, 600, 120);
1832        assert_eq!(
1833            anchor.timestamp_for(100),
1834            read_at - time::Duration::seconds(120)
1835        );
1836        assert_eq!(
1837            anchor.timestamp_for(99),
1838            read_at - time::Duration::seconds(120 + 600)
1839        );
1840    }
1841
1842    #[test]
1843    fn test_build_history_records_timestamps_follow_device_indices() {
1844        // Downloading only the three OLDEST records of a 1000-record buffer must
1845        // stamp them ~999 intervals ago, not as the newest three.
1846        let (anchor, read_at) = anchor_fixture(1000, 600, 120);
1847        let records = build_history_records(
1848            &anchor,
1849            1,
1850            &[800, 810, 820],
1851            &[450; 3],
1852            &[10130; 3],
1853            &[45; 3],
1854            &[],
1855        );
1856        let newest = read_at - time::Duration::seconds(120);
1857        assert_eq!(records.len(), 3);
1858        assert_eq!(
1859            records[0].timestamp,
1860            newest - time::Duration::seconds(999 * 600)
1861        );
1862        assert_eq!(
1863            records[2].timestamp,
1864            newest - time::Duration::seconds(997 * 600)
1865        );
1866    }
1867
1868    #[test]
1869    fn test_build_history_records_newest_record_matches_anchor() {
1870        let (anchor, read_at) = anchor_fixture(1000, 600, 120);
1871        let records = build_history_records(
1872            &anchor,
1873            998,
1874            &[800, 810, 820],
1875            &[450; 3],
1876            &[10130; 3],
1877            &[45; 3],
1878            &[],
1879        );
1880        assert_eq!(records[2].timestamp, read_at - time::Duration::seconds(120));
1881    }
1882}