fits_io/
slice_bin_table_hdu.rs1use crate::bin_table::BinTable;
2#[cfg(feature = "tokio")]
3use crate::bin_table::{FieldDefinition, OwnedRow};
4use crate::hdu::{BinTableHDU, HDU};
5use crate::header::Header;
6#[cfg(feature = "tokio")]
7use futures::StreamExt;
8#[cfg(feature = "tokio")]
9use futures::stream;
10#[cfg(feature = "tokio")]
11use futures::stream::BoxStream;
12#[cfg(feature = "serde")]
13use serde::de::DeserializeOwned;
14use std::error::Error;
15use std::sync::Arc;
16
17#[derive(Debug, Clone)]
19pub struct SliceBinTableHDU {
20 header: Header,
21 data: Arc<Vec<u8>>,
22 data_offset: usize,
23 pending: Option<Vec<u8>>,
26}
27
28impl SliceBinTableHDU {
29 pub(crate) fn new(header: Header, data: Arc<Vec<u8>>, data_offset: usize) -> Self {
30 Self {
31 header,
32 data,
33 data_offset,
34 pending: None,
35 }
36 }
37
38 pub fn from_table(table: &BinTable) -> Result<Self, Box<dyn Error + Send + Sync>> {
40 let mut hdu = Self {
41 header: Header::default(),
42 data: Arc::new(Vec::new()),
43 data_offset: 0,
44 pending: None,
45 };
46 hdu.set_table(table)?;
47
48 Ok(hdu)
49 }
50
51 pub(crate) fn data_bytes(&self) -> &[u8] {
53 if let Some(pending) = &self.pending {
54 return pending;
55 }
56
57 let len = self.header.data_bytes_len();
58 self.data
59 .get(self.data_offset..)
60 .and_then(|data| data.get(..len))
61 .unwrap_or_default()
62 }
63}
64
65impl HDU for SliceBinTableHDU {
66 fn header(&self) -> &Header {
67 &self.header
68 }
69
70 fn header_mut(&mut self) -> &mut Header {
71 &mut self.header
72 }
73}
74
75impl BinTableHDU for SliceBinTableHDU {
76 fn table_data_bytes_len(&self) -> u64 {
77 let (Some(bytes_per_row), Some(rows)) = (self.header.naxis_n(0), self.header.naxis_n(1))
78 else {
79 return 0;
80 };
81
82 (bytes_per_row.max(0) as u64).saturating_mul(rows.max(0) as u64)
83 }
84
85 fn read_table(&self) -> Result<BinTable, Box<dyn Error + Send + Sync>> {
86 BinTable::from_u8(&self.header, self.data_bytes().to_vec())
87 }
88
89 fn set_table(&mut self, table: &BinTable) -> Result<(), Box<dyn Error + Send + Sync>> {
90 crate::bin_table::write::apply_to_header(table, &mut self.header);
91 self.pending = Some(table.data().to_vec());
92
93 Ok(())
94 }
95
96 #[cfg(feature = "serde")]
97 fn read_rows<T: DeserializeOwned + Send + Sync>(
98 &self,
99 ) -> Result<Vec<T>, Box<dyn Error + Send + Sync>> {
100 let table = self.read_table()?;
101 Ok(crate::bin_table::from_bin_table(&table)?)
102 }
103
104 #[cfg(feature = "tokio")]
105 fn stream_table_rows_raw(
106 &self,
107 ) -> Result<BoxStream<'_, OwnedRow>, Box<dyn Error + Send + Sync>> {
108 let bytes_per_row = self.header.naxis_n(0).unwrap_or(0);
109 let bytes_per_row = usize::try_from(bytes_per_row)
110 .map_err(|_| format!("NAXIS1 must not be negative, but was {}", bytes_per_row))?;
111
112 if bytes_per_row == 0 {
113 return Ok(stream::empty().boxed());
114 }
115
116 let field_definitions = Arc::new(FieldDefinition::all_from_header(&self.header)?);
117
118 let data = self.data_bytes();
119 let heap_offset = self.header.table_heap_offset().min(data.len());
120 let heap = Arc::new(data[heap_offset..].to_vec());
121
122 let rows: Vec<_> = data[..heap_offset]
125 .chunks_exact(bytes_per_row)
126 .map(|row| {
127 OwnedRow::new(
128 Arc::clone(&field_definitions),
129 Arc::clone(&heap),
130 row.to_vec(),
131 )
132 })
133 .collect();
134
135 Ok(stream::iter(rows).boxed())
136 }
137
138 #[cfg(feature = "serde")]
139 #[cfg(feature = "tokio")]
140 fn stream_table_rows<T: DeserializeOwned + Send + Sync>(
141 &self,
142 ) -> Result<BoxStream<'_, crate::Result<T>>, Box<dyn Error + Send + Sync>> {
143 Ok(self
144 .stream_table_rows_raw()?
145 .map(|row| crate::bin_table::from_bin_table_row(&row.row()))
146 .boxed())
147 }
148}