1use 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#[derive(Debug, Clone)]
59pub struct HistoryProgress {
60 pub current_param: HistoryParam,
62 pub param_index: usize,
64 pub total_params: usize,
66 pub values_downloaded: usize,
68 pub total_values: usize,
70 pub overall_progress: f32,
72}
73
74impl HistoryProgress {
75 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 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
110pub type ProgressCallback = Arc<dyn Fn(HistoryProgress) + Send + Sync>;
112
113pub type CheckpointCallback = Arc<dyn Fn(HistoryCheckpoint) + Send + Sync>;
115
116#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
121pub struct HistoryCheckpoint {
122 pub device_id: String,
124 pub current_param: HistoryParamCheckpoint,
126 pub resume_index: u16,
128 pub total_readings: u16,
130 pub completed_params: Vec<HistoryParamCheckpoint>,
132 pub created_at: time::OffsetDateTime,
134 pub downloaded_data: Option<PartialHistoryData>,
136}
137
138#[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#[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 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 pub fn is_valid(&self, current_total_readings: u16) -> bool {
209 self.total_readings == current_total_readings
212 }
213
214 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 => {} }
225 }
226 }
227
228 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#[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 Humidity2 = 5,
247 Radon = 10,
249}
250
251#[derive(Clone)]
284pub struct HistoryOptions {
285 pub start_index: Option<u16>,
287 pub end_index: Option<u16>,
289 pub read_delay: Duration,
291 pub progress_callback: Option<ProgressCallback>,
293 pub use_adaptive_delay: bool,
295 pub checkpoint_callback: Option<CheckpointCallback>,
298 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, }
327 }
328}
329
330impl HistoryOptions {
331 #[must_use]
333 pub fn new() -> Self {
334 Self::default()
335 }
336
337 #[must_use]
339 pub fn start_index(mut self, index: u16) -> Self {
340 self.start_index = Some(index);
341 self
342 }
343
344 #[must_use]
346 pub fn end_index(mut self, index: u16) -> Self {
347 self.end_index = Some(index);
348 self
349 }
350
351 #[must_use]
353 pub fn read_delay(mut self, delay: Duration) -> Self {
354 self.read_delay = delay;
355 self
356 }
357
358 #[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 pub fn report_progress(&self, progress: &HistoryProgress) {
370 if let Some(cb) = &self.progress_callback {
371 cb(progress.clone());
372 }
373 }
374
375 #[must_use]
384 pub fn adaptive_delay(mut self, enable: bool) -> Self {
385 self.use_adaptive_delay = enable;
386 self
387 }
388
389 #[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 #[must_use]
406 pub fn checkpoint_interval(mut self, interval: usize) -> Self {
407 self.checkpoint_interval = interval;
408 self
409 }
410
411 #[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 pub fn report_checkpoint(&self, checkpoint: &HistoryCheckpoint) {
422 if let Some(cb) = &self.checkpoint_callback {
423 cb(checkpoint.clone());
424 }
425 }
426
427 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#[derive(Debug, Clone)]
443pub struct HistoryInfo {
444 pub total_readings: u16,
446 pub interval_seconds: u16,
448 pub seconds_since_update: u16,
450}
451
452#[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 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 pub async fn get_history_info(&self) -> Result<HistoryInfo> {
486 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 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 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 pub async fn download_history(&self) -> Result<Vec<HistoryRecord>> {
521 self.download_history_with_options(HistoryOptions::default())
522 .await
523 }
524
525 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 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 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 let effective_delay = options.effective_read_delay(signal_quality);
601
602 match self.device_type() {
604 Some(DeviceType::AranetRadiation) => {
605 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 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 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 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 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 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 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 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 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 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 #[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 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 let response = self.read_characteristic(HISTORY_V2).await?;
1035
1036 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 sleep(read_delay).await;
1065 continue;
1066 }
1067 consecutive_wrong_param = 0;
1068
1069 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 if resp_count == 0 {
1080 debug!("Reached end of history (count=0)");
1081 break;
1082 }
1083
1084 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 Ok(values.into_values().collect())
1129 }
1130
1131 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 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 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 let (tx, mut rx) = mpsc::channel::<Vec<u8>>(256);
1228
1229 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 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 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 let _ = self.unsubscribe_from_notifications(HISTORY_V1).await;
1267 return Err(e);
1268 }
1269
1270 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; 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 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 HistoryParam::Humidity2 | HistoryParam::Radon => {}
1340 }
1341 }
1342
1343 self.unsubscribe_from_notifications(HISTORY_V1).await?;
1345
1346 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 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#[derive(Debug, PartialEq, Eq)]
1393enum PacketProgress {
1394 Continue(u16),
1396 Done,
1398 Stalled,
1401}
1402
1403#[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 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 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 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
1453fn 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 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 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
1532pub fn raw_to_temperature(raw: u16) -> f32 {
1534 raw as f32 / 20.0
1535}
1536
1537pub fn raw_to_pressure(raw: u16) -> f32 {
1539 raw as f32 / 10.0
1540}
1541
1542#[cfg(test)]
1546mod tests {
1547 use super::*;
1548
1549 #[test]
1552 fn test_raw_to_temperature_typical_values() {
1553 assert!((raw_to_temperature(450) - 22.5).abs() < 0.001);
1555
1556 assert!((raw_to_temperature(400) - 20.0).abs() < 0.001);
1558
1559 assert!((raw_to_temperature(500) - 25.0).abs() < 0.001);
1561 }
1562
1563 #[test]
1564 fn test_raw_to_temperature_edge_cases() {
1565 assert!((raw_to_temperature(0) - 0.0).abs() < 0.001);
1567
1568 assert!((raw_to_temperature(1000) - 50.0).abs() < 0.001);
1573
1574 assert!((raw_to_temperature(u16::MAX) - 3276.75).abs() < 0.01);
1576 }
1577
1578 #[test]
1579 fn test_raw_to_temperature_precision() {
1580 assert!((raw_to_temperature(451) - 22.55).abs() < 0.001);
1583
1584 assert!((raw_to_temperature(441) - 22.05).abs() < 0.001);
1586 }
1587
1588 #[test]
1591 fn test_raw_to_pressure_typical_values() {
1592 assert!((raw_to_pressure(10132) - 1013.2).abs() < 0.01);
1594
1595 assert!((raw_to_pressure(10000) - 1000.0).abs() < 0.01);
1597
1598 assert!((raw_to_pressure(10500) - 1050.0).abs() < 0.01);
1600 }
1601
1602 #[test]
1603 fn test_raw_to_pressure_edge_cases() {
1604 assert!((raw_to_pressure(0) - 0.0).abs() < 0.01);
1606
1607 assert!((raw_to_pressure(9500) - 950.0).abs() < 0.01);
1609
1610 assert!((raw_to_pressure(11000) - 1100.0).abs() < 0.01);
1612
1613 assert!((raw_to_pressure(u16::MAX) - 6553.5).abs() < 0.1);
1615 }
1616
1617 #[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 #[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 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 #[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 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 let progress = apply_v2_packet(&mut values, 50, 100, 1, 3, &payload, 2, &parse_u16);
1755 assert_eq!(progress, PacketProgress::Stalled);
1756 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 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 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 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 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}