Skip to main content

runmat_runtime/data/
mod.rs

1use runmat_value::{ComplexStorage, ComplexTensor, IntegerComplexStorage};
2use std::cell::RefCell;
3use std::collections::{BTreeMap, HashMap};
4use std::future::Future;
5use std::path::{Path, PathBuf};
6use std::sync::atomic::{AtomicU64, Ordering};
7
8use chrono::Utc;
9use runmat_filesystem as fs;
10use runmat_filesystem::data_contract::{
11    DataChunkDescriptor, DataChunkUploadRequest, DataChunkUploadTarget,
12};
13use runmat_value::{
14    IntValue, IntegerStorage, NumericScalar, NumericStorage, ObjectInstance, Tensor, Value,
15};
16use serde::{Deserialize, Serialize};
17use sha2::{Digest, Sha256};
18
19use crate::builtins::math::elementwise::integer_cast::IntegerTarget;
20use crate::{build_runtime_error, BuiltinResult, RuntimeError};
21
22#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct DataManifest {
24    pub schema_version: u32,
25    pub format: String,
26    pub dataset_id: String,
27    pub name: Option<String>,
28    pub created_at: String,
29    pub updated_at: String,
30    pub arrays: BTreeMap<String, DataArrayMeta>,
31    pub attrs: BTreeMap<String, serde_json::Value>,
32    pub txn_sequence: u64,
33}
34
35#[derive(Debug, Clone, Serialize, Deserialize)]
36pub struct DataArrayMeta {
37    pub dtype: String,
38    pub shape: Vec<usize>,
39    pub chunk_shape: Vec<usize>,
40    #[serde(default = "default_array_order")]
41    pub order: String,
42    pub codec: String,
43    #[serde(default)]
44    pub chunk_index_path: Option<String>,
45    pub data_path: String,
46}
47
48fn default_array_order() -> String {
49    "column_major".to_string()
50}
51
52#[derive(Debug, Clone, Serialize, Deserialize)]
53pub struct DataArrayPayload {
54    pub dtype: String,
55    pub shape: Vec<usize>,
56    pub values: DataArrayValues,
57    #[serde(default, skip_serializing_if = "Option::is_none")]
58    pub imaginary_values: Option<DataArrayValues>,
59}
60
61/// The persisted backing values of a data-array payload.
62///
63/// JSON arrays written before integer storage was introduced are decoded as
64/// `F64`; new writes use the tagged representation below so every integer
65/// class can round-trip without passing through a floating point value.
66#[derive(Debug, Clone, PartialEq)]
67pub enum DataArrayValues {
68    F64(Vec<f64>),
69    F32(Vec<f32>),
70    I8(Vec<i8>),
71    I16(Vec<i16>),
72    I32(Vec<i32>),
73    I64(Vec<i64>),
74    U8(Vec<u8>),
75    U16(Vec<u16>),
76    U32(Vec<u32>),
77    U64(Vec<u64>),
78}
79
80#[derive(Serialize, Deserialize)]
81#[serde(tag = "encoding", content = "data", rename_all = "snake_case")]
82enum TaggedDataArrayValues {
83    F64(Vec<f64>),
84    F32(Vec<f32>),
85    I8(Vec<i8>),
86    I16(Vec<i16>),
87    I32(Vec<i32>),
88    I64(Vec<i64>),
89    U8(Vec<u8>),
90    U16(Vec<u16>),
91    U32(Vec<u32>),
92    U64(Vec<u64>),
93}
94
95#[derive(Serialize)]
96#[serde(tag = "encoding", content = "data", rename_all = "snake_case")]
97enum TaggedDataArrayValuesRef<'a> {
98    F64(&'a [f64]),
99    F32(&'a [f32]),
100    I8(&'a [i8]),
101    I16(&'a [i16]),
102    I32(&'a [i32]),
103    I64(&'a [i64]),
104    U8(&'a [u8]),
105    U16(&'a [u16]),
106    U32(&'a [u32]),
107    U64(&'a [u64]),
108}
109
110#[derive(Deserialize)]
111#[serde(untagged)]
112enum DataArrayValuesWire {
113    Tagged(TaggedDataArrayValues),
114    Legacy(Vec<f64>),
115}
116
117impl Serialize for DataArrayValues {
118    fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
119    where
120        S: serde::Serializer,
121    {
122        let tagged = match self {
123            Self::F64(values) => TaggedDataArrayValuesRef::F64(values),
124            Self::F32(values) => TaggedDataArrayValuesRef::F32(values),
125            Self::I8(values) => TaggedDataArrayValuesRef::I8(values),
126            Self::I16(values) => TaggedDataArrayValuesRef::I16(values),
127            Self::I32(values) => TaggedDataArrayValuesRef::I32(values),
128            Self::I64(values) => TaggedDataArrayValuesRef::I64(values),
129            Self::U8(values) => TaggedDataArrayValuesRef::U8(values),
130            Self::U16(values) => TaggedDataArrayValuesRef::U16(values),
131            Self::U32(values) => TaggedDataArrayValuesRef::U32(values),
132            Self::U64(values) => TaggedDataArrayValuesRef::U64(values),
133        };
134        tagged.serialize(serializer)
135    }
136}
137
138impl<'de> Deserialize<'de> for DataArrayValues {
139    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
140    where
141        D: serde::Deserializer<'de>,
142    {
143        Ok(match DataArrayValuesWire::deserialize(deserializer)? {
144            DataArrayValuesWire::Legacy(values) => Self::F64(values),
145            DataArrayValuesWire::Tagged(tagged) => match tagged {
146                TaggedDataArrayValues::F64(values) => Self::F64(values),
147                TaggedDataArrayValues::F32(values) => Self::F32(values),
148                TaggedDataArrayValues::I8(values) => Self::I8(values),
149                TaggedDataArrayValues::I16(values) => Self::I16(values),
150                TaggedDataArrayValues::I32(values) => Self::I32(values),
151                TaggedDataArrayValues::I64(values) => Self::I64(values),
152                TaggedDataArrayValues::U8(values) => Self::U8(values),
153                TaggedDataArrayValues::U16(values) => Self::U16(values),
154                TaggedDataArrayValues::U32(values) => Self::U32(values),
155                TaggedDataArrayValues::U64(values) => Self::U64(values),
156            },
157        })
158    }
159}
160
161impl DataArrayValues {
162    pub fn zeros(dtype: &str, len: usize) -> BuiltinResult<Self> {
163        if is_single_dtype(dtype) {
164            return Ok(Self::F32(vec![0.0; len]));
165        }
166        Ok(match integer_dtype(dtype) {
167            Some("int8") => Self::I8(vec![0; len]),
168            Some("int16") => Self::I16(vec![0; len]),
169            Some("int32") => Self::I32(vec![0; len]),
170            Some("int64") => Self::I64(vec![0; len]),
171            Some("uint8") => Self::U8(vec![0; len]),
172            Some("uint16") => Self::U16(vec![0; len]),
173            Some("uint32") => Self::U32(vec![0; len]),
174            Some("uint64") => Self::U64(vec![0; len]),
175            None if is_double_dtype(dtype) => Self::F64(vec![0.0; len]),
176            _ => return Err(unsupported_data_dtype(dtype)),
177        })
178    }
179
180    pub fn len(&self) -> usize {
181        match self {
182            Self::F64(values) => values.len(),
183            Self::F32(values) => values.len(),
184            Self::I8(values) => values.len(),
185            Self::I16(values) => values.len(),
186            Self::I32(values) => values.len(),
187            Self::I64(values) => values.len(),
188            Self::U8(values) => values.len(),
189            Self::U16(values) => values.len(),
190            Self::U32(values) => values.len(),
191            Self::U64(values) => values.len(),
192        }
193    }
194
195    pub fn is_empty(&self) -> bool {
196        self.len() == 0
197    }
198
199    pub fn into_tensor(self, shape: Vec<usize>) -> Result<Tensor, String> {
200        match self {
201            Self::F64(values) => Tensor::new(values, shape),
202            Self::F32(values) => Tensor::from_f32(values, shape),
203            Self::I8(values) => Tensor::new_integer(IntegerStorage::I8(values), shape),
204            Self::I16(values) => Tensor::new_integer(IntegerStorage::I16(values), shape),
205            Self::I32(values) => Tensor::new_integer(IntegerStorage::I32(values), shape),
206            Self::I64(values) => Tensor::new_integer(IntegerStorage::I64(values), shape),
207            Self::U8(values) => Tensor::new_integer(IntegerStorage::U8(values), shape),
208            Self::U16(values) => Tensor::new_integer(IntegerStorage::U16(values), shape),
209            Self::U32(values) => Tensor::new_integer(IntegerStorage::U32(values), shape),
210            Self::U64(values) => Tensor::new_integer(IntegerStorage::U64(values), shape),
211        }
212    }
213
214    pub fn to_f64_vec(&self) -> Vec<f64> {
215        match self {
216            Self::F64(values) => values.clone(),
217            Self::F32(values) => values.iter().map(|&value| f64::from(value)).collect(),
218            Self::I8(values) => values.iter().map(|&value| value as f64).collect(),
219            Self::I16(values) => values.iter().map(|&value| value as f64).collect(),
220            Self::I32(values) => values.iter().map(|&value| value as f64).collect(),
221            Self::I64(values) => values.iter().map(|&value| value as f64).collect(),
222            Self::U8(values) => values.iter().map(|&value| value as f64).collect(),
223            Self::U16(values) => values.iter().map(|&value| value as f64).collect(),
224            Self::U32(values) => values.iter().map(|&value| value as f64).collect(),
225            Self::U64(values) => values.iter().map(|&value| value as f64).collect(),
226        }
227    }
228
229    /// Returns at most `limit` values converted to the numeric preview format.
230    ///
231    /// Data-file previews cross the WASM boundary as JavaScript numbers. Keep
232    /// that established representation while only converting the requested
233    /// prefix: a fallback full-payload read may contain substantially more
234    /// values than the preview is allowed to return.
235    pub fn preview_f64(&self, limit: usize) -> Vec<f64> {
236        match self {
237            Self::F64(values) => values.iter().take(limit).copied().collect(),
238            Self::F32(values) => values
239                .iter()
240                .take(limit)
241                .map(|&value| f64::from(value))
242                .collect(),
243            Self::I8(values) => values
244                .iter()
245                .take(limit)
246                .map(|&value| value as f64)
247                .collect(),
248            Self::I16(values) => values
249                .iter()
250                .take(limit)
251                .map(|&value| value as f64)
252                .collect(),
253            Self::I32(values) => values
254                .iter()
255                .take(limit)
256                .map(|&value| value as f64)
257                .collect(),
258            Self::I64(values) => values
259                .iter()
260                .take(limit)
261                .map(|&value| value as f64)
262                .collect(),
263            Self::U8(values) => values
264                .iter()
265                .take(limit)
266                .map(|&value| value as f64)
267                .collect(),
268            Self::U16(values) => values
269                .iter()
270                .take(limit)
271                .map(|&value| value as f64)
272                .collect(),
273            Self::U32(values) => values
274                .iter()
275                .take(limit)
276                .map(|&value| value as f64)
277                .collect(),
278            Self::U64(values) => values
279                .iter()
280                .take(limit)
281                .map(|&value| value as f64)
282                .collect(),
283        }
284    }
285
286    pub fn get(&self, index: usize) -> BuiltinResult<DataScalar> {
287        match self {
288            Self::F64(values) => values.get(index).copied().map(DataScalar::F64),
289            Self::F32(values) => values.get(index).copied().map(DataScalar::F32),
290            Self::I8(values) => values.get(index).copied().map(|v| DataScalar::I8(v)),
291            Self::I16(values) => values.get(index).copied().map(|v| DataScalar::I16(v)),
292            Self::I32(values) => values.get(index).copied().map(|v| DataScalar::I32(v)),
293            Self::I64(values) => values.get(index).copied().map(|v| DataScalar::I64(v)),
294            Self::U8(values) => values.get(index).copied().map(|v| DataScalar::U8(v)),
295            Self::U16(values) => values.get(index).copied().map(|v| DataScalar::U16(v)),
296            Self::U32(values) => values.get(index).copied().map(|v| DataScalar::U32(v)),
297            Self::U64(values) => values.get(index).copied().map(|v| DataScalar::U64(v)),
298        }
299        .ok_or_else(|| data_error(format!("data payload index {index} is out of bounds")))
300    }
301
302    pub fn push(&mut self, value: DataScalar) -> BuiltinResult<()> {
303        match (self, value) {
304            (Self::F64(values), DataScalar::F64(value)) => values.push(value),
305            (Self::F32(values), DataScalar::F32(value)) => values.push(value),
306            (Self::I8(values), DataScalar::I8(value)) => values.push(value),
307            (Self::I16(values), DataScalar::I16(value)) => values.push(value),
308            (Self::I32(values), DataScalar::I32(value)) => values.push(value),
309            (Self::I64(values), DataScalar::I64(value)) => values.push(value),
310            (Self::U8(values), DataScalar::U8(value)) => values.push(value),
311            (Self::U16(values), DataScalar::U16(value)) => values.push(value),
312            (Self::U32(values), DataScalar::U32(value)) => values.push(value),
313            (Self::U64(values), DataScalar::U64(value)) => values.push(value),
314            _ => return Err(data_error("data payload storage class mismatch")),
315        }
316        Ok(())
317    }
318
319    pub fn set(&mut self, index: usize, value: DataScalar) -> BuiltinResult<()> {
320        match (self, value) {
321            (Self::F64(values), DataScalar::F64(value)) => set_at(values, index, value),
322            (Self::F32(values), DataScalar::F32(value)) => set_at(values, index, value),
323            (Self::I8(values), DataScalar::I8(value)) => set_at(values, index, value),
324            (Self::I16(values), DataScalar::I16(value)) => set_at(values, index, value),
325            (Self::I32(values), DataScalar::I32(value)) => set_at(values, index, value),
326            (Self::I64(values), DataScalar::I64(value)) => set_at(values, index, value),
327            (Self::U8(values), DataScalar::U8(value)) => set_at(values, index, value),
328            (Self::U16(values), DataScalar::U16(value)) => set_at(values, index, value),
329            (Self::U32(values), DataScalar::U32(value)) => set_at(values, index, value),
330            (Self::U64(values), DataScalar::U64(value)) => set_at(values, index, value),
331            _ => return Err(data_error("data payload storage class mismatch")),
332        }?;
333        Ok(())
334    }
335
336    fn cast_to_dtype(self, dtype: &str) -> BuiltinResult<Self> {
337        if is_single_dtype(dtype) {
338            return Ok(Self::F32(
339                self.to_f64_vec()
340                    .into_iter()
341                    .map(|value| value as f32)
342                    .collect(),
343            ));
344        }
345        if is_double_dtype(dtype) {
346            return Ok(Self::F64(self.to_f64_vec()));
347        }
348        let Some(target) = integer_target(dtype) else {
349            return Err(unsupported_data_dtype(dtype));
350        };
351        let mut values = Vec::with_capacity(self.len());
352        for index in 0..self.len() {
353            let value = self.get(index)?;
354            values.push(match value {
355                DataScalar::F64(value) => target.cast_scalar(value),
356                DataScalar::F32(value) => target.cast_scalar(f64::from(value)),
357                value => target.cast_int(&value.to_int_value()),
358            });
359        }
360        Ok(Self::from_integer_storage(target.storage(values)))
361    }
362
363    fn from_integer_storage(storage: IntegerStorage) -> Self {
364        match storage {
365            IntegerStorage::I8(values) => Self::I8(values),
366            IntegerStorage::I16(values) => Self::I16(values),
367            IntegerStorage::I32(values) => Self::I32(values),
368            IntegerStorage::I64(values) => Self::I64(values),
369            IntegerStorage::U8(values) => Self::U8(values),
370            IntegerStorage::U16(values) => Self::U16(values),
371            IntegerStorage::U32(values) => Self::U32(values),
372            IntegerStorage::U64(values) => Self::U64(values),
373        }
374    }
375
376    fn from_numeric_storage(storage: NumericStorage) -> Self {
377        match storage {
378            NumericStorage::F64(values) => Self::F64(values),
379            NumericStorage::F32(values) => Self::F32(values),
380            NumericStorage::I8(values) => Self::I8(values),
381            NumericStorage::I16(values) => Self::I16(values),
382            NumericStorage::I32(values) => Self::I32(values),
383            NumericStorage::I64(values) => Self::I64(values),
384            NumericStorage::U8(values) => Self::U8(values),
385            NumericStorage::U16(values) => Self::U16(values),
386            NumericStorage::U32(values) => Self::U32(values),
387            NumericStorage::U64(values) => Self::U64(values),
388        }
389    }
390
391    fn into_numeric_storage(self) -> NumericStorage {
392        match self {
393            Self::F64(values) => NumericStorage::F64(values),
394            Self::F32(values) => NumericStorage::F32(values),
395            Self::I8(values) => NumericStorage::I8(values),
396            Self::I16(values) => NumericStorage::I16(values),
397            Self::I32(values) => NumericStorage::I32(values),
398            Self::I64(values) => NumericStorage::I64(values),
399            Self::U8(values) => NumericStorage::U8(values),
400            Self::U16(values) => NumericStorage::U16(values),
401            Self::U32(values) => NumericStorage::U32(values),
402            Self::U64(values) => NumericStorage::U64(values),
403        }
404    }
405}
406
407#[derive(Debug, Clone, Copy)]
408pub enum DataScalar {
409    F64(f64),
410    F32(f32),
411    I8(i8),
412    I16(i16),
413    I32(i32),
414    I64(i64),
415    U8(u8),
416    U16(u16),
417    U32(u32),
418    U64(u64),
419}
420
421impl DataScalar {
422    fn to_int_value(self) -> IntValue {
423        match self {
424            Self::F64(value) => IntValue::I64(value as i64),
425            Self::F32(value) => IntValue::I64(value as i64),
426            Self::I8(value) => IntValue::I8(value),
427            Self::I16(value) => IntValue::I16(value),
428            Self::I32(value) => IntValue::I32(value),
429            Self::I64(value) => IntValue::I64(value),
430            Self::U8(value) => IntValue::U8(value),
431            Self::U16(value) => IntValue::U16(value),
432            Self::U32(value) => IntValue::U32(value),
433            Self::U64(value) => IntValue::U64(value),
434        }
435    }
436}
437
438fn set_at<T>(values: &mut [T], index: usize, value: T) -> BuiltinResult<()> {
439    let target = values
440        .get_mut(index)
441        .ok_or_else(|| data_error(format!("data payload index {index} is out of bounds")))?;
442    *target = value;
443    Ok(())
444}
445
446fn integer_dtype(dtype: &str) -> Option<&'static str> {
447    match dtype.to_ascii_lowercase().as_str() {
448        "int8" => Some("int8"),
449        "int16" => Some("int16"),
450        "int32" => Some("int32"),
451        "int64" => Some("int64"),
452        "uint8" => Some("uint8"),
453        "uint16" => Some("uint16"),
454        "uint32" => Some("uint32"),
455        "uint64" => Some("uint64"),
456        _ => None,
457    }
458}
459
460fn is_single_dtype(dtype: &str) -> bool {
461    matches!(
462        dtype.to_ascii_lowercase().as_str(),
463        "single" | "f32" | "float32"
464    )
465}
466
467fn is_double_dtype(dtype: &str) -> bool {
468    matches!(
469        dtype.to_ascii_lowercase().as_str(),
470        "double" | "f64" | "float64"
471    )
472}
473
474fn unsupported_data_dtype(dtype: &str) -> RuntimeError {
475    data_error(format!(
476        "unsupported data array dtype '{dtype}'; expected f64, f32, or a built-in integer class"
477    ))
478}
479
480fn integer_target(dtype: &str) -> Option<IntegerTarget> {
481    match integer_dtype(dtype) {
482        Some("int8") => Some(IntegerTarget::I8),
483        Some("int16") => Some(IntegerTarget::I16),
484        Some("int32") => Some(IntegerTarget::I32),
485        Some("int64") => Some(IntegerTarget::I64),
486        Some("uint8") => Some(IntegerTarget::U8),
487        Some("uint16") => Some(IntegerTarget::U16),
488        Some("uint32") => Some(IntegerTarget::U32),
489        Some("uint64") => Some(IntegerTarget::U64),
490        _ => None,
491    }
492}
493
494impl DataArrayPayload {
495    pub fn zeros(dtype: String, shape: Vec<usize>) -> BuiltinResult<Self> {
496        let values = DataArrayValues::zeros(&dtype, checked_shape_element_count(&shape)?)?;
497        Ok(Self {
498            dtype,
499            shape,
500            values,
501            imaginary_values: None,
502        })
503    }
504
505    pub fn from_value(dtype: String, value: &Value) -> BuiltinResult<Self> {
506        let (shape, values, imaginary_values) = data_values_from_value(value)?;
507        Ok(Self {
508            dtype: dtype.clone(),
509            shape,
510            values: values.cast_to_dtype(&dtype)?,
511            imaginary_values: imaginary_values
512                .map(|values| values.cast_to_dtype(&dtype))
513                .transpose()?,
514        })
515    }
516
517    pub fn filled(dtype: String, shape: Vec<usize>, value: &Value) -> BuiltinResult<Self> {
518        let scalar = Self::from_value(dtype.clone(), value)?;
519        if scalar.values.len() != 1 {
520            return Err(data_error("expected numeric scalar"));
521        }
522        let imaginary_scalar = scalar
523            .imaginary_values
524            .as_ref()
525            .map(|values| values.get(0))
526            .transpose()?;
527        let scalar = scalar.values.get(0)?;
528        let len = checked_shape_element_count(&shape)?;
529        let mut values = DataArrayValues::zeros(&dtype, len)?;
530        let mut imaginary_values = imaginary_scalar
531            .map(|_| DataArrayValues::zeros(&dtype, len))
532            .transpose()?;
533        for index in 0..len {
534            values.set(index, scalar)?;
535            if let (Some(values), Some(scalar)) = (&mut imaginary_values, imaginary_scalar) {
536                values.set(index, scalar)?;
537            }
538        }
539        Ok(Self {
540            dtype,
541            shape,
542            values,
543            imaginary_values,
544        })
545    }
546
547    pub fn normalize_for_dtype(mut self, dtype: &str) -> BuiltinResult<Self> {
548        self.values = self.values.cast_to_dtype(dtype)?;
549        self.imaginary_values = self
550            .imaginary_values
551            .map(|values| values.cast_to_dtype(dtype))
552            .transpose()?;
553        self.dtype = dtype.to_string();
554        Ok(self)
555    }
556
557    pub fn into_value(self) -> BuiltinResult<Value> {
558        let Some(imaginary_values) = self.imaginary_values else {
559            return self
560                .values
561                .into_tensor(self.shape)
562                .map(Value::Tensor)
563                .map_err(|err| data_error(format!("invalid data payload: {err}")));
564        };
565        let real = self.values.into_numeric_storage();
566        let imag = imaginary_values.into_numeric_storage();
567        let storage = match (real, imag) {
568            (NumericStorage::F64(real), NumericStorage::F64(imag)) => {
569                ComplexStorage::F64(real.into_iter().zip(imag).collect())
570            }
571            (NumericStorage::F32(real), NumericStorage::F32(imag)) => {
572                ComplexStorage::F32(real.into_iter().zip(imag).collect())
573            }
574            (real, imag) => {
575                let real = real.into_integer_storage().map_err(|_| {
576                    data_error("complex data payload components have mismatched storage classes")
577                })?;
578                let imag = imag.into_integer_storage().map_err(|_| {
579                    data_error("complex data payload components have mismatched storage classes")
580                })?;
581                ComplexStorage::Integer(IntegerComplexStorage::new(real, imag).map_err(
582                    |error| data_error(format!("invalid complex data payload: {error}")),
583                )?)
584            }
585        };
586        ComplexTensor::from_complex_storage(storage, self.shape)
587            .map(Value::ComplexTensor)
588            .map_err(|error| data_error(format!("invalid complex data payload: {error}")))
589    }
590}
591
592fn checked_shape_element_count(shape: &[usize]) -> BuiltinResult<usize> {
593    if shape.contains(&0) {
594        return Ok(0);
595    }
596    shape.iter().try_fold(1usize, |count, dimension| {
597        count
598            .checked_mul(*dimension)
599            .ok_or_else(|| data_error("data array shape exceeds platform element-count limits"))
600    })
601}
602
603fn data_values_from_value(
604    value: &Value,
605) -> BuiltinResult<(Vec<usize>, DataArrayValues, Option<DataArrayValues>)> {
606    match value {
607        Value::Tensor(tensor) => {
608            let storage = tensor
609                .clone()
610                .into_numeric_storage()
611                .map_err(|error| data_error(format!("invalid numeric tensor storage: {error}")))?;
612            let values = DataArrayValues::from_numeric_storage(storage);
613            Ok((tensor.shape.clone(), values, None))
614        }
615        Value::Num(value) => Ok((vec![1, 1], DataArrayValues::F64(vec![*value]), None)),
616        Value::Int(IntValue::I8(value)) => {
617            Ok((vec![1, 1], DataArrayValues::I8(vec![*value]), None))
618        }
619        Value::Int(IntValue::I16(value)) => {
620            Ok((vec![1, 1], DataArrayValues::I16(vec![*value]), None))
621        }
622        Value::Int(IntValue::I32(value)) => {
623            Ok((vec![1, 1], DataArrayValues::I32(vec![*value]), None))
624        }
625        Value::Int(IntValue::I64(value)) => {
626            Ok((vec![1, 1], DataArrayValues::I64(vec![*value]), None))
627        }
628        Value::Int(IntValue::U8(value)) => {
629            Ok((vec![1, 1], DataArrayValues::U8(vec![*value]), None))
630        }
631        Value::Int(IntValue::U16(value)) => {
632            Ok((vec![1, 1], DataArrayValues::U16(vec![*value]), None))
633        }
634        Value::Int(IntValue::U32(value)) => {
635            Ok((vec![1, 1], DataArrayValues::U32(vec![*value]), None))
636        }
637        Value::Int(IntValue::U64(value)) => {
638            Ok((vec![1, 1], DataArrayValues::U64(vec![*value]), None))
639        }
640        Value::Complex(real, imag) => Ok((
641            vec![1, 1],
642            DataArrayValues::F64(vec![*real]),
643            Some(DataArrayValues::F64(vec![*imag])),
644        )),
645        Value::ComplexTensor(tensor) => {
646            let shape = tensor.shape.clone();
647            let (real, imag) = match tensor.clone().into_complex_storage() {
648                ComplexStorage::F64(values) => {
649                    let (real, imag): (Vec<_>, Vec<_>) = values.into_iter().unzip();
650                    (DataArrayValues::F64(real), DataArrayValues::F64(imag))
651                }
652                ComplexStorage::F32(values) => {
653                    let (real, imag): (Vec<_>, Vec<_>) = values.into_iter().unzip();
654                    (DataArrayValues::F32(real), DataArrayValues::F32(imag))
655                }
656                ComplexStorage::Integer(storage) => (
657                    DataArrayValues::from_integer_storage(storage.real),
658                    DataArrayValues::from_integer_storage(storage.imag),
659                ),
660            };
661            Ok((shape, real, Some(imag)))
662        }
663        _ => Err(data_error(
664            "DataArray.write supports tensor or numeric scalar values",
665        )),
666    }
667}
668
669#[derive(Debug, Clone, Serialize, Deserialize)]
670pub struct DataChunkIndex {
671    pub schema_version: u32,
672    pub array: String,
673    pub chunks: Vec<DataChunkIndexEntry>,
674}
675
676#[derive(Debug, Clone, Serialize, Deserialize)]
677pub struct DataChunkIndexEntry {
678    pub key: String,
679    pub object_id: String,
680    pub hash: String,
681    pub bytes_raw: u64,
682    pub bytes_stored: u64,
683    #[serde(default)]
684    pub coords: Vec<usize>,
685    #[serde(default)]
686    pub shape: Vec<usize>,
687    pub data_path: String,
688}
689
690#[derive(Debug, Clone)]
691pub struct DataSchema {
692    pub arrays: BTreeMap<String, DataArrayMeta>,
693}
694
695#[derive(Debug, Clone)]
696pub struct PendingTxn {
697    pub dataset_path: String,
698    pub base_sequence: u64,
699    pub writes: Vec<PendingWrite>,
700    pub resizes: Vec<PendingResize>,
701    pub fills: Vec<PendingFill>,
702    pub create_arrays: Vec<PendingCreateArray>,
703    pub delete_arrays: Vec<String>,
704    pub attrs: BTreeMap<String, Value>,
705    pub status: TxnStatus,
706}
707
708#[derive(Debug, Clone)]
709pub struct PendingWrite {
710    pub array: String,
711    pub slice_spec: Option<Value>,
712    pub value: Value,
713}
714
715#[derive(Debug, Clone)]
716pub struct PendingResize {
717    pub array: String,
718    pub shape: Vec<usize>,
719}
720
721#[derive(Debug, Clone)]
722pub struct PendingFill {
723    pub array: String,
724    pub slice_spec: Option<Value>,
725    pub value: Value,
726}
727
728#[derive(Debug, Clone)]
729pub struct PendingCreateArray {
730    pub array: String,
731    pub meta: DataArrayMeta,
732}
733
734#[derive(Debug, Clone, PartialEq, Eq)]
735pub enum TxnStatus {
736    Open,
737    Committed,
738    Aborted,
739}
740
741thread_local! {
742    static FALLBACK_TX_REGISTRY: RefCell<HashMap<String, PendingTxn>> = RefCell::new(HashMap::new());
743}
744
745#[cfg(not(target_arch = "wasm32"))]
746tokio::task_local! {
747    static TASK_TX_REGISTRY: RefCell<HashMap<String, PendingTxn>>;
748}
749
750pub async fn with_tx_registry_scope<F>(future: F) -> F::Output
751where
752    F: Future,
753{
754    #[cfg(not(target_arch = "wasm32"))]
755    {
756        if TASK_TX_REGISTRY.try_with(|_| ()).is_ok() {
757            future.await
758        } else {
759            TASK_TX_REGISTRY
760                .scope(RefCell::new(HashMap::new()), future)
761                .await
762        }
763    }
764    #[cfg(target_arch = "wasm32")]
765    {
766        future.await
767    }
768}
769
770fn with_tx_registry<T>(f: impl FnOnce(&mut HashMap<String, PendingTxn>) -> T) -> BuiltinResult<T> {
771    #[cfg(not(target_arch = "wasm32"))]
772    {
773        if TASK_TX_REGISTRY.try_with(|_| ()).is_ok() {
774            return TASK_TX_REGISTRY.with(|registry| {
775                let mut registry = registry.try_borrow_mut().map_err(|_| {
776                    data_error("data transaction registry is already mutably borrowed")
777                })?;
778                Ok(f(&mut registry))
779            });
780        }
781    }
782
783    FALLBACK_TX_REGISTRY.with(|registry| {
784        let mut registry = registry
785            .try_borrow_mut()
786            .map_err(|_| data_error("data transaction registry is already mutably borrowed"))?;
787        Ok(f(&mut registry))
788    })
789}
790
791pub fn data_error(message: impl Into<String>) -> RuntimeError {
792    build_runtime_error(message)
793        .with_identifier("RUNMAT:Data:Error")
794        .with_builtin("data")
795        .build()
796}
797
798fn data_error_with_identifier(
799    message: impl Into<String>,
800    identifier: &'static str,
801) -> RuntimeError {
802    build_runtime_error(message)
803        .with_identifier(identifier)
804        .with_builtin("data")
805        .build()
806}
807
808const DATA_MANIFEST_CONFLICT_IDENTIFIER: &str = "RunMat:data:ManifestConflict";
809const DATA_TRANSACTION_NOT_FOUND_IDENTIFIER: &str = "RunMat:data:TransactionNotFound";
810
811pub fn parse_string(value: &Value, context: &str) -> BuiltinResult<String> {
812    match value {
813        Value::String(s) => Ok(s.clone()),
814        Value::CharArray(chars) => chars
815            .row_string()
816            .ok_or_else(|| data_error(format!("{context}: expected character row vector"))),
817        _ => Err(data_error(format!("{context}: expected string value"))),
818    }
819}
820
821pub fn dataset_root(path: &str) -> PathBuf {
822    PathBuf::from(path)
823}
824
825pub fn manifest_path(root: &Path) -> PathBuf {
826    root.join("manifest.json")
827}
828
829pub fn arrays_root(root: &Path) -> PathBuf {
830    root.join("arrays")
831}
832
833pub async fn write_manifest_async(root: &Path, manifest: &DataManifest) -> BuiltinResult<()> {
834    fs::create_dir_all_async(root).await.map_err(|err| {
835        data_error(format!(
836            "failed to create dataset root '{}': {err}",
837            root.display()
838        ))
839    })?;
840    let path = manifest_path(root);
841    let bytes = serde_json::to_vec_pretty(manifest)
842        .map_err(|err| data_error(format!("failed to encode manifest json: {err}")))?;
843    fs::write_async(&path, &bytes).await.map_err(|err| {
844        data_error(format!(
845            "failed to write manifest '{}': {err}",
846            path.display()
847        ))
848    })?;
849    Ok(())
850}
851
852pub async fn read_manifest_async(root: &Path) -> BuiltinResult<DataManifest> {
853    let path = manifest_path(root);
854    let bytes = fs::read_async(&path).await.map_err(|err| {
855        data_error(format!(
856            "failed to read manifest '{}': {err}",
857            path.display()
858        ))
859    })?;
860    let manifest = serde_json::from_slice::<DataManifest>(&bytes).map_err(|err| {
861        data_error(format!(
862            "failed to parse manifest '{}': {err}",
863            path.display()
864        ))
865    })?;
866    Ok(manifest)
867}
868
869pub async fn write_array_payload_async(
870    root: &Path,
871    array: &str,
872    payload: &DataArrayPayload,
873    chunk_shape: &[usize],
874) -> BuiltinResult<(PathBuf, PathBuf)> {
875    let array_dir = arrays_root(root).join(array);
876    fs::create_dir_all_async(&array_dir).await.map_err(|err| {
877        data_error(format!(
878            "failed to create array dir '{}': {err}",
879            array_dir.display()
880        ))
881    })?;
882    let payload_path = array_dir.join("data.f64.json");
883    let bytes = serde_json::to_vec(payload)
884        .map_err(|err| data_error(format!("failed to encode array payload json: {err}")))?;
885    fs::write_async(&payload_path, &bytes)
886        .await
887        .map_err(|err| {
888            data_error(format!(
889                "failed to write payload '{}': {err}",
890                payload_path.display()
891            ))
892        })?;
893
894    let chunk_dir = array_dir.join("chunks");
895    fs::create_dir_all_async(&chunk_dir).await.map_err(|err| {
896        data_error(format!(
897            "failed to create chunk dir '{}': {err}",
898            chunk_dir.display()
899        ))
900    })?;
901
902    let mut index = DataChunkIndex {
903        schema_version: 1,
904        array: array.to_string(),
905        chunks: Vec::new(),
906    };
907    let mut upload_chunks = Vec::new();
908    let grid_shape = chunk_grid_shape(&payload.shape, chunk_shape);
909    let mut coords = vec![0usize; payload.shape.len()];
910    loop {
911        let chunk_start = chunk_start_for_coords(&coords, chunk_shape);
912        let chunk_extent = chunk_extent_for_start(&chunk_start, chunk_shape, &payload.shape);
913        let chunk_payload = DataArrayPayload {
914            dtype: payload.dtype.clone(),
915            shape: chunk_extent.clone(),
916            values: collect_chunk_values(payload, &chunk_start, &chunk_extent)?,
917            imaginary_values: payload
918                .imaginary_values
919                .as_ref()
920                .map(|values| {
921                    collect_chunk_component_values(
922                        values,
923                        &payload.dtype,
924                        &payload.shape,
925                        &chunk_start,
926                        &chunk_extent,
927                    )
928                })
929                .transpose()?,
930        };
931        let key = chunk_key(&coords);
932        let object_id = format!("obj_{}", key.replace('.', "_"));
933        let chunk_bytes = serde_json::to_vec(&chunk_payload)
934            .map_err(|err| data_error(format!("failed to encode chunk payload: {err}")))?;
935        let data_path = chunk_dir.join(format!("{object_id}.json"));
936        fs::write_async(&data_path, &chunk_bytes)
937            .await
938            .map_err(|err| {
939                data_error(format!(
940                    "failed to write chunk '{}': {err}",
941                    data_path.display()
942                ))
943            })?;
944        let hash = sha256_hex(&chunk_bytes);
945        let rel_chunk_path = data_path
946            .strip_prefix(root)
947            .map_err(|err| data_error(format!("failed to compute chunk relative path: {err}")))?
948            .to_string_lossy()
949            .to_string();
950        index.chunks.push(DataChunkIndexEntry {
951            key: key.clone(),
952            object_id: object_id.clone(),
953            hash: hash.clone(),
954            bytes_raw: chunk_bytes.len() as u64,
955            bytes_stored: chunk_bytes.len() as u64,
956            coords: coords.clone(),
957            shape: chunk_extent,
958            data_path: rel_chunk_path,
959        });
960        upload_chunks.push((
961            DataChunkDescriptor {
962                key,
963                object_id,
964                hash,
965                bytes_raw: chunk_bytes.len() as u64,
966                bytes_stored: chunk_bytes.len() as u64,
967            },
968            chunk_bytes,
969        ));
970        if !advance_index(&mut coords, &grid_shape) {
971            break;
972        }
973    }
974
975    maybe_upload_chunks_async(root, array, upload_chunks).await?;
976
977    tracing::info!(
978        target: "runmat.data",
979        dataset = %root.display(),
980        array = array,
981        chunks = index.chunks.len(),
982        payload_bytes = bytes.len(),
983        "data chunk write planned"
984    );
985
986    let chunk_index_path = chunk_dir.join("index.json");
987    let chunk_index_bytes = serde_json::to_vec(&index)
988        .map_err(|err| data_error(format!("failed to encode chunk index json: {err}")))?;
989    fs::write_async(&chunk_index_path, &chunk_index_bytes)
990        .await
991        .map_err(|err| {
992            data_error(format!(
993                "failed to write chunk index '{}': {err}",
994                chunk_index_path.display()
995            ))
996        })?;
997    Ok((payload_path, chunk_index_path))
998}
999
1000pub async fn read_array_payload_async(
1001    root: &Path,
1002    meta: &DataArrayMeta,
1003) -> BuiltinResult<DataArrayPayload> {
1004    if let Some(index_path) = &meta.chunk_index_path {
1005        let path = root.join(index_path);
1006        if fs::metadata_async(&path).await.is_ok() {
1007            return read_array_payload_chunked_async(root, meta, &path).await;
1008        }
1009    }
1010    let payload_path = root.join(&meta.data_path);
1011    let bytes = fs::read_async(&payload_path).await.map_err(|err| {
1012        data_error(format!(
1013            "failed to read payload '{}': {err}",
1014            payload_path.display()
1015        ))
1016    })?;
1017    serde_json::from_slice::<DataArrayPayload>(&bytes)
1018        .map_err(|err| {
1019            data_error(format!(
1020                "failed to parse payload '{}': {err}",
1021                payload_path.display()
1022            ))
1023        })?
1024        .normalize_for_dtype(&meta.dtype)
1025}
1026
1027pub async fn read_array_slice_payload_async(
1028    root: &Path,
1029    meta: &DataArrayMeta,
1030    start: &[usize],
1031    shape: &[usize],
1032) -> BuiltinResult<DataArrayPayload> {
1033    let (slice_start, slice_shape) = normalize_slice_bounds(&meta.shape, start, shape)?;
1034    if let Some(index_path) = &meta.chunk_index_path {
1035        let path = root.join(index_path);
1036        if fs::metadata_async(&path).await.is_ok() {
1037            return read_array_payload_chunked_slice_async(
1038                root,
1039                meta,
1040                &path,
1041                &slice_start,
1042                &slice_shape,
1043            )
1044            .await;
1045        }
1046    }
1047    let full = read_array_payload_async(root, meta).await?;
1048    extract_slice_payload(&full, &slice_start, &slice_shape)
1049}
1050
1051async fn read_array_payload_chunked_slice_async(
1052    root: &Path,
1053    meta: &DataArrayMeta,
1054    index_path: &Path,
1055    slice_start: &[usize],
1056    slice_shape: &[usize],
1057) -> BuiltinResult<DataArrayPayload> {
1058    let bytes = fs::read_async(index_path).await.map_err(|err| {
1059        data_error(format!(
1060            "failed to read chunk index '{}': {err}",
1061            index_path.display()
1062        ))
1063    })?;
1064    let index: DataChunkIndex = serde_json::from_slice(&bytes).map_err(|err| {
1065        data_error(format!(
1066            "failed to parse chunk index '{}': {err}",
1067            index_path.display()
1068        ))
1069    })?;
1070
1071    let mut values =
1072        DataArrayValues::zeros(&meta.dtype, checked_shape_element_count(slice_shape)?)?;
1073    let mut imaginary_values: Option<DataArrayValues> = None;
1074    for chunk in index.chunks {
1075        let coords = chunk_coords_from_entry(&chunk, meta.shape.len())?;
1076        let chunk_start = chunk_start_for_coords(&coords, &meta.chunk_shape);
1077        let chunk_extent = if chunk.shape.is_empty() {
1078            chunk_extent_for_start(&chunk_start, &meta.chunk_shape, &meta.shape)
1079        } else {
1080            chunk.shape.clone()
1081        };
1082        if !chunk_intersects_slice(&chunk_start, &chunk_extent, slice_start, slice_shape) {
1083            continue;
1084        }
1085
1086        let chunk_path = root.join(&chunk.data_path);
1087        let bytes = fs::read_async(&chunk_path).await.map_err(|err| {
1088            data_error(format!(
1089                "failed to read chunk payload '{}': {err}",
1090                chunk_path.display()
1091            ))
1092        })?;
1093        let payload: DataArrayPayload = serde_json::from_slice::<DataArrayPayload>(&bytes)
1094            .map_err(|err| {
1095                data_error(format!(
1096                    "failed to parse chunk payload '{}': {err}",
1097                    chunk_path.display()
1098                ))
1099            })?
1100            .normalize_for_dtype(&meta.dtype)?;
1101        if payload.shape != chunk_extent {
1102            return Err(data_error(format!(
1103                "chunk payload shape mismatch for key '{}': {:?} != {:?}",
1104                chunk.key, payload.shape, chunk_extent
1105            )));
1106        }
1107
1108        let mut local = vec![0usize; chunk_extent.len()];
1109        loop {
1110            let mut global = Vec::with_capacity(chunk_extent.len());
1111            for dim in 0..chunk_extent.len() {
1112                global.push(chunk_start[dim] + local[dim]);
1113            }
1114            if coordinate_in_slice(&global, slice_start, slice_shape) {
1115                let src_linear = linear_index_column_major(&local, &chunk_extent)?;
1116                let mut dst = Vec::with_capacity(slice_shape.len());
1117                for dim in 0..slice_shape.len() {
1118                    dst.push(global[dim].saturating_sub(slice_start[dim]));
1119                }
1120                let dst_linear = linear_index_column_major(&dst, slice_shape)?;
1121                values.set(dst_linear, payload.values.get(src_linear)?)?;
1122                if let Some(payload_imaginary) = &payload.imaginary_values {
1123                    let target = match &mut imaginary_values {
1124                        Some(values) => values,
1125                        None => imaginary_values.insert(DataArrayValues::zeros(
1126                            &meta.dtype,
1127                            checked_shape_element_count(slice_shape)?,
1128                        )?),
1129                    };
1130                    target.set(dst_linear, payload_imaginary.get(src_linear)?)?;
1131                }
1132            }
1133            if !advance_index(&mut local, &chunk_extent) {
1134                break;
1135            }
1136        }
1137    }
1138
1139    Ok(DataArrayPayload {
1140        dtype: meta.dtype.clone(),
1141        shape: slice_shape.to_vec(),
1142        values,
1143        imaginary_values,
1144    })
1145}
1146
1147async fn read_array_payload_chunked_async(
1148    root: &Path,
1149    meta: &DataArrayMeta,
1150    index_path: &Path,
1151) -> BuiltinResult<DataArrayPayload> {
1152    let bytes = fs::read_async(index_path).await.map_err(|err| {
1153        data_error(format!(
1154            "failed to read chunk index '{}': {err}",
1155            index_path.display()
1156        ))
1157    })?;
1158    let index: DataChunkIndex = serde_json::from_slice(&bytes).map_err(|err| {
1159        data_error(format!(
1160            "failed to parse chunk index '{}': {err}",
1161            index_path.display()
1162        ))
1163    })?;
1164    let mut values =
1165        DataArrayValues::zeros(&meta.dtype, checked_shape_element_count(&meta.shape)?)?;
1166    let mut imaginary_values: Option<DataArrayValues> = None;
1167    for chunk in index.chunks {
1168        let chunk_path = root.join(&chunk.data_path);
1169        let bytes = fs::read_async(&chunk_path).await.map_err(|err| {
1170            data_error(format!(
1171                "failed to read chunk payload '{}': {err}",
1172                chunk_path.display()
1173            ))
1174        })?;
1175        let payload: DataArrayPayload = serde_json::from_slice::<DataArrayPayload>(&bytes)
1176            .map_err(|err| {
1177                data_error(format!(
1178                    "failed to parse chunk payload '{}': {err}",
1179                    chunk_path.display()
1180                ))
1181            })?
1182            .normalize_for_dtype(&meta.dtype)?;
1183        let coords = chunk_coords_from_entry(&chunk, meta.shape.len())?;
1184        let chunk_start = chunk_start_for_coords(&coords, &meta.chunk_shape);
1185        let chunk_extent = if chunk.shape.is_empty() {
1186            chunk_extent_for_start(&chunk_start, &meta.chunk_shape, &meta.shape)
1187        } else {
1188            chunk.shape.clone()
1189        };
1190        if payload.shape != chunk_extent {
1191            return Err(data_error(format!(
1192                "chunk payload shape mismatch for key '{}': {:?} != {:?}",
1193                chunk.key, payload.shape, chunk_extent
1194            )));
1195        }
1196        let mut local = vec![0usize; chunk_extent.len()];
1197        loop {
1198            let mut global = Vec::with_capacity(chunk_extent.len());
1199            for dim in 0..chunk_extent.len() {
1200                global.push(chunk_start[dim] + local[dim]);
1201            }
1202            let src_linear = linear_index_column_major(&local, &chunk_extent)?;
1203            let dst_linear = linear_index_column_major(&global, &meta.shape)?;
1204            values.set(dst_linear, payload.values.get(src_linear)?)?;
1205            if let Some(payload_imaginary) = &payload.imaginary_values {
1206                let target = match &mut imaginary_values {
1207                    Some(values) => values,
1208                    None => imaginary_values.insert(DataArrayValues::zeros(
1209                        &meta.dtype,
1210                        checked_shape_element_count(&meta.shape)?,
1211                    )?),
1212                };
1213                target.set(dst_linear, payload_imaginary.get(src_linear)?)?;
1214            }
1215            if !advance_index(&mut local, &chunk_extent) {
1216                break;
1217            }
1218        }
1219    }
1220    Ok(DataArrayPayload {
1221        dtype: meta.dtype.clone(),
1222        shape: meta.shape.clone(),
1223        values,
1224        imaginary_values,
1225    })
1226}
1227
1228async fn maybe_upload_chunks_async(
1229    root: &Path,
1230    array: &str,
1231    chunks: Vec<(DataChunkDescriptor, Vec<u8>)>,
1232) -> BuiltinResult<()> {
1233    if chunks.is_empty() {
1234        return Ok(());
1235    }
1236    let request = DataChunkUploadRequest {
1237        dataset_path: root.to_string_lossy().to_string(),
1238        array: array.to_string(),
1239        chunks: chunks.iter().map(|(desc, _)| desc.clone()).collect(),
1240    };
1241    let targets = match fs::data_chunk_upload_targets_async(&request).await {
1242        Ok(targets) => targets,
1243        Err(err) if err.kind() == std::io::ErrorKind::Unsupported => return Ok(()),
1244        Err(err) => {
1245            return Err(data_error(format!(
1246                "failed to request data chunk upload targets: {err}"
1247            )))
1248        }
1249    };
1250    for (descriptor, bytes) in chunks {
1251        let target = find_chunk_target(&targets, &descriptor.key)?;
1252        fs::data_upload_chunk_async(target, &bytes)
1253            .await
1254            .map_err(|err| {
1255                data_error(format!(
1256                    "failed to upload chunk '{}': {err}",
1257                    descriptor.key
1258                ))
1259            })?;
1260        tracing::info!(
1261            target: "runmat.data",
1262            dataset = %root.display(),
1263            array = array,
1264            chunk_key = descriptor.key,
1265            bytes = bytes.len(),
1266            "data chunk uploaded"
1267        );
1268    }
1269    Ok(())
1270}
1271
1272fn find_chunk_target<'a>(
1273    targets: &'a [DataChunkUploadTarget],
1274    key: &str,
1275) -> BuiltinResult<&'a DataChunkUploadTarget> {
1276    targets
1277        .iter()
1278        .find(|target| target.key == key)
1279        .ok_or_else(|| data_error(format!("missing upload target for chunk '{key}'")))
1280}
1281
1282pub fn sha256_hex(bytes: &[u8]) -> String {
1283    let mut hasher = Sha256::new();
1284    hasher.update(bytes);
1285    let digest = hasher.finalize();
1286    format!("sha256:{:x}", digest)
1287}
1288
1289fn chunk_key(coords: &[usize]) -> String {
1290    coords
1291        .iter()
1292        .map(|v| v.to_string())
1293        .collect::<Vec<_>>()
1294        .join(".")
1295}
1296
1297fn chunk_grid_shape(shape: &[usize], chunk_shape: &[usize]) -> Vec<usize> {
1298    shape
1299        .iter()
1300        .enumerate()
1301        .map(|(idx, extent)| {
1302            let chunk = chunk_shape.get(idx).copied().unwrap_or(1).max(1);
1303            extent.div_ceil(chunk)
1304        })
1305        .collect()
1306}
1307
1308fn chunk_start_for_coords(coords: &[usize], chunk_shape: &[usize]) -> Vec<usize> {
1309    coords
1310        .iter()
1311        .enumerate()
1312        .map(|(idx, coord)| coord * chunk_shape.get(idx).copied().unwrap_or(1).max(1))
1313        .collect()
1314}
1315
1316fn chunk_extent_for_start(
1317    start: &[usize],
1318    chunk_shape: &[usize],
1319    full_shape: &[usize],
1320) -> Vec<usize> {
1321    start
1322        .iter()
1323        .enumerate()
1324        .map(|(idx, start)| {
1325            let chunk = chunk_shape.get(idx).copied().unwrap_or(1).max(1);
1326            let end = (*start + chunk).min(full_shape[idx]);
1327            end.saturating_sub(*start)
1328        })
1329        .collect()
1330}
1331
1332fn collect_chunk_values(
1333    payload: &DataArrayPayload,
1334    chunk_start: &[usize],
1335    chunk_extent: &[usize],
1336) -> BuiltinResult<DataArrayValues> {
1337    collect_chunk_component_values(
1338        &payload.values,
1339        &payload.dtype,
1340        &payload.shape,
1341        chunk_start,
1342        chunk_extent,
1343    )
1344}
1345
1346fn collect_chunk_component_values(
1347    source: &DataArrayValues,
1348    dtype: &str,
1349    full_shape: &[usize],
1350    chunk_start: &[usize],
1351    chunk_extent: &[usize],
1352) -> BuiltinResult<DataArrayValues> {
1353    let mut local = vec![0usize; chunk_extent.len()];
1354    let mut values = DataArrayValues::zeros(dtype, 0)?;
1355    loop {
1356        let mut global = Vec::with_capacity(chunk_extent.len());
1357        for dim in 0..chunk_extent.len() {
1358            global.push(chunk_start[dim] + local[dim]);
1359        }
1360        let linear = linear_index_column_major(&global, full_shape)?;
1361        values.push(source.get(linear)?)?;
1362        if !advance_index(&mut local, chunk_extent) {
1363            break;
1364        }
1365    }
1366    Ok(values)
1367}
1368
1369fn chunk_coords_from_entry(entry: &DataChunkIndexEntry, rank: usize) -> BuiltinResult<Vec<usize>> {
1370    if !entry.coords.is_empty() {
1371        if entry.coords.len() != rank {
1372            return Err(data_error(format!(
1373                "chunk coords rank mismatch for key '{}': expected {rank}, got {}",
1374                entry.key,
1375                entry.coords.len()
1376            )));
1377        }
1378        return Ok(entry.coords.clone());
1379    }
1380    let coords = entry
1381        .key
1382        .split('.')
1383        .map(|part| {
1384            part.parse::<usize>()
1385                .map_err(|_| data_error(format!("invalid chunk key '{}'", entry.key)))
1386        })
1387        .collect::<BuiltinResult<Vec<_>>>()?;
1388    if coords.len() != rank {
1389        return Err(data_error(format!(
1390            "chunk key rank mismatch for key '{}': expected {rank}, got {}",
1391            entry.key,
1392            coords.len()
1393        )));
1394    }
1395    Ok(coords)
1396}
1397
1398fn normalize_slice_bounds(
1399    full_shape: &[usize],
1400    start: &[usize],
1401    shape: &[usize],
1402) -> BuiltinResult<(Vec<usize>, Vec<usize>)> {
1403    if full_shape.is_empty() {
1404        return Ok((Vec::new(), Vec::new()));
1405    }
1406    let mut normalized_start = Vec::with_capacity(full_shape.len());
1407    let mut normalized_shape = Vec::with_capacity(full_shape.len());
1408    for (axis, axis_len) in full_shape.iter().copied().enumerate() {
1409        if axis_len == 0 {
1410            return Err(data_error("slice axis length must be greater than zero"));
1411        }
1412        let requested_start = start.get(axis).copied().unwrap_or(0);
1413        let clamped_start = requested_start.min(axis_len.saturating_sub(1));
1414        let requested_span = shape.get(axis).copied().unwrap_or(axis_len);
1415        let clamped_span = requested_span
1416            .max(1)
1417            .min(axis_len.saturating_sub(clamped_start));
1418        normalized_start.push(clamped_start);
1419        normalized_shape.push(clamped_span);
1420    }
1421    Ok((normalized_start, normalized_shape))
1422}
1423
1424fn coordinate_in_slice(global: &[usize], slice_start: &[usize], slice_shape: &[usize]) -> bool {
1425    for dim in 0..slice_shape.len() {
1426        let start = slice_start[dim];
1427        let end = start.saturating_add(slice_shape[dim]);
1428        let value = global[dim];
1429        if value < start || value >= end {
1430            return false;
1431        }
1432    }
1433    true
1434}
1435
1436fn chunk_intersects_slice(
1437    chunk_start: &[usize],
1438    chunk_extent: &[usize],
1439    slice_start: &[usize],
1440    slice_shape: &[usize],
1441) -> bool {
1442    for dim in 0..slice_shape.len() {
1443        let chunk_lo = chunk_start[dim];
1444        let chunk_hi = chunk_lo.saturating_add(chunk_extent[dim]);
1445        let slice_lo = slice_start[dim];
1446        let slice_hi = slice_lo.saturating_add(slice_shape[dim]);
1447        if chunk_hi <= slice_lo || slice_hi <= chunk_lo {
1448            return false;
1449        }
1450    }
1451    true
1452}
1453
1454fn extract_slice_payload(
1455    payload: &DataArrayPayload,
1456    start: &[usize],
1457    shape: &[usize],
1458) -> BuiltinResult<DataArrayPayload> {
1459    let mut values = DataArrayValues::zeros(&payload.dtype, 0)?;
1460    let mut imaginary_values = payload
1461        .imaginary_values
1462        .as_ref()
1463        .map(|_| DataArrayValues::zeros(&payload.dtype, 0))
1464        .transpose()?;
1465    if shape.is_empty() {
1466        return Ok(DataArrayPayload {
1467            dtype: payload.dtype.clone(),
1468            shape: Vec::new(),
1469            values,
1470            imaginary_values,
1471        });
1472    }
1473    let mut local = vec![0usize; shape.len()];
1474    loop {
1475        let mut global = Vec::with_capacity(shape.len());
1476        for dim in 0..shape.len() {
1477            global.push(start[dim] + local[dim]);
1478        }
1479        let linear = linear_index_column_major(&global, &payload.shape)?;
1480        values.push(payload.values.get(linear)?)?;
1481        if let (Some(source), Some(target)) = (&payload.imaginary_values, &mut imaginary_values) {
1482            target.push(source.get(linear)?)?;
1483        }
1484        if !advance_index(&mut local, shape) {
1485            break;
1486        }
1487    }
1488    Ok(DataArrayPayload {
1489        dtype: payload.dtype.clone(),
1490        shape: shape.to_vec(),
1491        values,
1492        imaginary_values,
1493    })
1494}
1495
1496fn linear_index_column_major(index: &[usize], shape: &[usize]) -> BuiltinResult<usize> {
1497    if index.len() != shape.len() {
1498        return Err(data_error("chunk index rank mismatch"));
1499    }
1500    let mut stride = 1usize;
1501    let mut linear = 0usize;
1502    for (idx, extent) in index.iter().zip(shape.iter()) {
1503        if *idx >= *extent {
1504            return Err(data_error("chunk index out of bounds"));
1505        }
1506        linear += idx * stride;
1507        stride = stride.saturating_mul(*extent);
1508    }
1509    Ok(linear)
1510}
1511
1512fn advance_index(index: &mut [usize], shape: &[usize]) -> bool {
1513    if shape.is_empty() {
1514        return false;
1515    }
1516    for dim in 0..shape.len() {
1517        index[dim] += 1;
1518        if index[dim] < shape[dim] {
1519            return true;
1520        }
1521        index[dim] = 0;
1522    }
1523    false
1524}
1525
1526pub fn parse_schema(schema: &Value) -> BuiltinResult<DataSchema> {
1527    let Value::Struct(schema_struct) = schema else {
1528        return Err(data_error("data.create: schema must be a struct"));
1529    };
1530    let arrays_value = schema_struct
1531        .fields
1532        .get("arrays")
1533        .ok_or_else(|| data_error("data.create: schema missing 'arrays' field"))?;
1534    let Value::Struct(arrays_struct) = arrays_value else {
1535        return Err(data_error("data.create: schema.arrays must be a struct"));
1536    };
1537
1538    let mut arrays = BTreeMap::new();
1539    for (name, meta_value) in &arrays_struct.fields {
1540        let Value::Struct(meta_struct) = meta_value else {
1541            return Err(data_error(format!(
1542                "data.create: schema.arrays.{name} must be a struct"
1543            )));
1544        };
1545        let dtype = meta_struct
1546            .fields
1547            .get("dtype")
1548            .map(|v| parse_string(v, "data.create schema dtype"))
1549            .transpose()?
1550            .unwrap_or_else(|| "f64".to_string());
1551        let shape = meta_struct
1552            .fields
1553            .get("shape")
1554            .map(parse_usize_vector)
1555            .transpose()?
1556            .unwrap_or_else(|| vec![0, 0]);
1557        let chunk_shape = meta_struct
1558            .fields
1559            .get("chunk")
1560            .map(parse_usize_vector)
1561            .transpose()?
1562            .unwrap_or_else(|| default_chunk_shape(&shape));
1563        validate_chunk_shape(&shape, &chunk_shape)?;
1564        let codec = meta_struct
1565            .fields
1566            .get("codec")
1567            .map(|v| parse_string(v, "data.create schema codec"))
1568            .transpose()?
1569            .unwrap_or_else(|| "zstd".to_string());
1570        let data_path = format!("arrays/{name}/data.f64.json");
1571        let chunk_index_path = format!("arrays/{name}/chunks/index.json");
1572        arrays.insert(
1573            name.clone(),
1574            DataArrayMeta {
1575                dtype,
1576                shape,
1577                chunk_shape,
1578                order: default_array_order(),
1579                codec,
1580                chunk_index_path: Some(chunk_index_path),
1581                data_path,
1582            },
1583        );
1584    }
1585
1586    Ok(DataSchema { arrays })
1587}
1588
1589fn default_chunk_shape(shape: &[usize]) -> Vec<usize> {
1590    if shape.is_empty() {
1591        return Vec::new();
1592    }
1593    let mut out = shape.to_vec();
1594    if out.len() == 1 {
1595        out[0] = out[0].clamp(1, 65_536);
1596        return out;
1597    }
1598    out[0] = out[0].clamp(1, 256);
1599    out[1] = out[1].clamp(1, 256);
1600    for dim in out.iter_mut().skip(2) {
1601        *dim = (*dim).clamp(1, 8);
1602    }
1603    out
1604}
1605
1606pub fn validate_chunk_shape(shape: &[usize], chunk_shape: &[usize]) -> BuiltinResult<()> {
1607    if chunk_shape.len() != shape.len() {
1608        return Err(data_error(
1609            "data array chunk shape must have the same rank as its array shape",
1610        ));
1611    }
1612    if chunk_shape.contains(&0) {
1613        return Err(data_error(
1614            "data array chunk dimensions must be strictly positive",
1615        ));
1616    }
1617    Ok(())
1618}
1619
1620fn parse_usize_vector(value: &Value) -> BuiltinResult<Vec<usize>> {
1621    match value {
1622        Value::Tensor(t) => tensor_to_usize_vector(t),
1623        Value::Num(n) => floating_dimension_to_usize(*n).map(|value| vec![value]),
1624        Value::Int(i) => i
1625            .try_to_usize()
1626            .map(|n| vec![n])
1627            .ok_or_else(|| data_error("data schema dimensions must be non-negative integers")),
1628        _ => Err(data_error(
1629            "data schema dimension field must be numeric tensor/vector",
1630        )),
1631    }
1632}
1633
1634fn floating_dimension_to_usize(value: f64) -> BuiltinResult<usize> {
1635    if !value.is_finite() || value < 0.0 || value.fract() != 0.0 {
1636        return Err(data_error(
1637            "data schema dimensions must be non-negative finite integers",
1638        ));
1639    }
1640    let integer = value as u128;
1641    if integer > usize::MAX as u128 || integer as f64 != value {
1642        return Err(data_error("data schema dimensions exceed platform limits"));
1643    }
1644    usize::try_from(integer)
1645        .map_err(|_| data_error("data schema dimensions exceed platform limits"))
1646}
1647
1648fn tensor_to_usize_vector(t: &Tensor) -> BuiltinResult<Vec<usize>> {
1649    let mut out = Vec::with_capacity(t.len());
1650    for index in 0..t.len() {
1651        let value = t
1652            .numeric_value_at(index)
1653            .ok_or_else(|| data_error("data schema dimensions require valid numeric storage"))?;
1654        out.push(match value {
1655            NumericScalar::F64(value) => floating_dimension_to_usize(value)?,
1656            NumericScalar::F32(value) => floating_dimension_to_usize(f64::from(value))?,
1657            value => value
1658                .into_int_value()
1659                .and_then(|value| value.try_to_usize())
1660                .ok_or_else(|| {
1661                    data_error("data schema dimensions must be non-negative integers")
1662                })?,
1663        });
1664    }
1665    Ok(out)
1666}
1667
1668pub fn dataset_object(path: &str, manifest: &DataManifest) -> Value {
1669    let mut obj = ObjectInstance::new("Dataset".to_string());
1670    obj.properties
1671        .insert("__data_path".to_string(), Value::String(path.to_string()));
1672    obj.properties.insert(
1673        "__data_id".to_string(),
1674        Value::String(manifest.dataset_id.clone()),
1675    );
1676    obj.properties.insert(
1677        "__data_version".to_string(),
1678        Value::String(manifest_version_token(manifest)),
1679    );
1680    Value::Object(obj)
1681}
1682
1683pub fn manifest_version_token(manifest: &DataManifest) -> String {
1684    format!("{}:{}", manifest.updated_at, manifest.txn_sequence)
1685}
1686
1687pub fn ensure_manifest_sequence(expected: u64, manifest: &DataManifest) -> BuiltinResult<()> {
1688    if manifest.txn_sequence != expected {
1689        tracing::warn!(
1690            target: "runmat.data",
1691            expected_sequence = expected,
1692            actual_sequence = manifest.txn_sequence,
1693            "manifest conflict detected"
1694        );
1695        return Err(data_error_with_identifier(
1696            "MANIFEST_CONFLICT: dataset changed since transaction begin",
1697            DATA_MANIFEST_CONFLICT_IDENTIFIER,
1698        ));
1699    }
1700    Ok(())
1701}
1702
1703pub fn array_object(dataset_path: &str, array_name: &str) -> Value {
1704    let mut obj = ObjectInstance::new("DataArray".to_string());
1705    obj.properties.insert(
1706        "__data_path".to_string(),
1707        Value::String(dataset_path.to_string()),
1708    );
1709    obj.properties.insert(
1710        "__array_name".to_string(),
1711        Value::String(array_name.to_string()),
1712    );
1713    Value::Object(obj)
1714}
1715
1716pub fn transaction_object(dataset_path: &str, tx_id: &str) -> Value {
1717    let mut obj = ObjectInstance::new("DataTransaction".to_string());
1718    obj.properties.insert(
1719        "__data_path".to_string(),
1720        Value::String(dataset_path.to_string()),
1721    );
1722    obj.properties
1723        .insert("__tx_id".to_string(), Value::String(tx_id.to_string()));
1724    Value::Object(obj)
1725}
1726
1727pub fn get_object_prop<'a>(obj: &'a ObjectInstance, key: &str) -> BuiltinResult<&'a Value> {
1728    obj.properties
1729        .get(key)
1730        .ok_or_else(|| data_error(format!("object missing internal property '{key}'")))
1731}
1732
1733pub fn now_rfc3339() -> String {
1734    Utc::now().to_rfc3339_opts(chrono::SecondsFormat::Secs, true)
1735}
1736
1737pub fn new_dataset_id() -> String {
1738    static NEXT_DATASET_ID: AtomicU64 = AtomicU64::new(1);
1739    let seq = NEXT_DATASET_ID.fetch_add(1, Ordering::Relaxed);
1740    format!("ds_{}_{}", Utc::now().timestamp_millis(), seq)
1741}
1742
1743pub fn new_tx_id() -> String {
1744    static NEXT_TX_ID: AtomicU64 = AtomicU64::new(1);
1745    let seq = NEXT_TX_ID.fetch_add(1, Ordering::Relaxed);
1746    format!("tx_{}_{}", Utc::now().timestamp_millis(), seq)
1747}
1748
1749pub fn start_tx(dataset_path: String, base_sequence: u64) -> BuiltinResult<String> {
1750    let tx_id = new_tx_id();
1751    let pending = PendingTxn {
1752        dataset_path,
1753        base_sequence,
1754        writes: Vec::new(),
1755        resizes: Vec::new(),
1756        fills: Vec::new(),
1757        create_arrays: Vec::new(),
1758        delete_arrays: Vec::new(),
1759        attrs: BTreeMap::new(),
1760        status: TxnStatus::Open,
1761    };
1762    with_tx_registry(|registry| {
1763        registry.insert(tx_id.clone(), pending);
1764    })?;
1765    Ok(tx_id)
1766}
1767
1768pub fn with_tx_mut<T>(
1769    tx_id: &str,
1770    f: impl FnOnce(&mut PendingTxn) -> BuiltinResult<T>,
1771) -> BuiltinResult<T> {
1772    with_tx_registry(|registry| {
1773        let tx = registry.get_mut(tx_id).ok_or_else(|| {
1774            data_error_with_identifier(
1775                format!("transaction '{tx_id}' not found"),
1776                DATA_TRANSACTION_NOT_FOUND_IDENTIFIER,
1777            )
1778        })?;
1779        f(tx)
1780    })?
1781}
1782
1783pub fn with_tx<T>(
1784    tx_id: &str,
1785    f: impl FnOnce(&PendingTxn) -> BuiltinResult<T>,
1786) -> BuiltinResult<T> {
1787    #[cfg(not(target_arch = "wasm32"))]
1788    {
1789        if TASK_TX_REGISTRY.try_with(|_| ()).is_ok() {
1790            return TASK_TX_REGISTRY.with(|registry| {
1791                let registry = registry
1792                    .try_borrow()
1793                    .map_err(|_| data_error("data transaction registry is already borrowed"))?;
1794                let tx = registry.get(tx_id).ok_or_else(|| {
1795                    data_error_with_identifier(
1796                        format!("transaction '{tx_id}' not found"),
1797                        DATA_TRANSACTION_NOT_FOUND_IDENTIFIER,
1798                    )
1799                })?;
1800                f(tx)
1801            });
1802        }
1803    }
1804
1805    FALLBACK_TX_REGISTRY.with(|registry| {
1806        let registry = registry
1807            .try_borrow()
1808            .map_err(|_| data_error("data transaction registry is already borrowed"))?;
1809        let tx = registry.get(tx_id).ok_or_else(|| {
1810            data_error_with_identifier(
1811                format!("transaction '{tx_id}' not found"),
1812                DATA_TRANSACTION_NOT_FOUND_IDENTIFIER,
1813            )
1814        })?;
1815        f(tx)
1816    })
1817}
1818
1819pub fn remove_tx(tx_id: &str) -> BuiltinResult<()> {
1820    with_tx_registry(|registry| {
1821        let _ = registry.remove(tx_id);
1822    })
1823}
1824
1825#[cfg(test)]
1826mod tests {
1827    use super::*;
1828
1829    #[test]
1830    fn schema_dimensions_preserve_large_typed_unsigned_values() {
1831        let expected = usize::try_from(u64::MAX).ok();
1832        let parsed = parse_usize_vector(&Value::Int(IntValue::U64(u64::MAX)));
1833        match expected {
1834            Some(value) => assert_eq!(parsed.expect("representable dimension"), vec![value]),
1835            None => assert!(parsed.is_err()),
1836        }
1837        assert!(parse_usize_vector(&Value::Int(IntValue::I64(-1))).is_err());
1838    }
1839
1840    #[test]
1841    fn schema_dimension_tensors_preserve_exact_integer_storage() {
1842        let cases = [
1843            IntegerStorage::I8(vec![2, 3]),
1844            IntegerStorage::I16(vec![2, 3]),
1845            IntegerStorage::I32(vec![2, 3]),
1846            IntegerStorage::I64(vec![2, 3]),
1847            IntegerStorage::U8(vec![2, 3]),
1848            IntegerStorage::U16(vec![2, 3]),
1849            IntegerStorage::U32(vec![2, 3]),
1850            IntegerStorage::U64(vec![2, 3]),
1851        ];
1852
1853        for storage in cases {
1854            let input = Tensor::new_integer(storage, vec![1, 2]).expect("dimension tensor");
1855            assert_eq!(
1856                parse_usize_vector(&Value::Tensor(input)).expect("typed dimensions"),
1857                vec![2, 3]
1858            );
1859        }
1860
1861        #[cfg(target_pointer_width = "64")]
1862        {
1863            let input = Tensor::new_integer(
1864                IntegerStorage::U64(vec![1_u64 << 53, (1_u64 << 53) + 1]),
1865                vec![1, 2],
1866            )
1867            .expect("uint64 dimension tensor");
1868            assert_eq!(
1869                parse_usize_vector(&Value::Tensor(input)).expect("wide typed dimensions"),
1870                vec![
1871                    usize::try_from(1_u64 << 53).expect("representable first dimension"),
1872                    usize::try_from((1_u64 << 53) + 1).expect("representable second dimension"),
1873                ]
1874            );
1875        }
1876    }
1877
1878    #[test]
1879    fn schema_dimension_tensors_accept_native_single_and_reject_fractional_floats() {
1880        let input = Tensor::from_f32(vec![2.0, 3.0], vec![1, 2]).expect("single dimensions");
1881        assert_eq!(
1882            parse_usize_vector(&Value::Tensor(input)).expect("single dimensions"),
1883            vec![2, 3]
1884        );
1885
1886        let fractional =
1887            Tensor::from_f32(vec![2.0, 3.5], vec![1, 2]).expect("fractional dimensions");
1888        assert!(parse_usize_vector(&Value::Tensor(fractional)).is_err());
1889        assert!(parse_usize_vector(&Value::Num(1.5)).is_err());
1890        assert!(parse_usize_vector(&Value::Num(f64::INFINITY)).is_err());
1891    }
1892
1893    #[test]
1894    fn schema_dimension_tensors_reject_negative_integer_storage() {
1895        let input =
1896            Tensor::new_integer(IntegerStorage::I16(vec![2, -1]), vec![1, 2]).expect("int16 dims");
1897        assert!(parse_usize_vector(&Value::Tensor(input)).is_err());
1898    }
1899
1900    #[test]
1901    fn payload_allocation_rejects_shape_product_overflow() {
1902        let error = DataArrayPayload::zeros("uint64".to_string(), vec![usize::MAX, 2])
1903            .expect_err("overflowing shape must reject before allocation");
1904        assert!(error
1905            .message()
1906            .contains("shape exceeds platform element-count limits"));
1907        let error = DataArrayPayload::filled(
1908            "uint64".to_string(),
1909            vec![usize::MAX, 2],
1910            &Value::Int(IntValue::U64(1)),
1911        )
1912        .expect_err("overflowing filled shape must reject before allocation");
1913        assert!(error
1914            .message()
1915            .contains("shape exceeds platform element-count limits"));
1916        for shape in [
1917            vec![0, usize::MAX, 2],
1918            vec![usize::MAX, 0, 2],
1919            vec![usize::MAX, 2, 0],
1920        ] {
1921            let payload = DataArrayPayload::zeros("uint64".to_string(), shape.clone())
1922                .expect("a zero dimension makes the total element count zero");
1923            assert_eq!(payload.shape, shape);
1924            assert_eq!(payload.values.len(), 0);
1925        }
1926    }
1927
1928    #[test]
1929    fn chunk_shape_requires_positive_rank_matched_dimensions() {
1930        validate_chunk_shape(&[4, 5], &[2, 5]).expect("valid chunk shape");
1931        assert!(validate_chunk_shape(&[4, 5], &[2]).is_err());
1932        assert!(validate_chunk_shape(&[4, 5], &[2, 0]).is_err());
1933        validate_chunk_shape(&[], &[]).expect("rank-zero metadata remains self-consistent");
1934    }
1935
1936    #[test]
1937    fn payload_rejects_unknown_dtype_instead_of_falling_back_to_f64() {
1938        let error = DataArrayPayload::zeros("mystery".to_string(), vec![1, 1])
1939            .expect_err("unknown dtype must reject");
1940        assert!(error.message().contains("unsupported data array dtype"));
1941        let error = DataArrayPayload::from_value("mystery".to_string(), &Value::Num(1.0))
1942            .expect_err("unknown cast target must reject");
1943        assert!(error.message().contains("unsupported data array dtype"));
1944    }
1945
1946    #[test]
1947    fn payload_roundtrips_every_native_integer_storage_class() {
1948        let cases = vec![
1949            DataArrayValues::I8(vec![i8::MIN, i8::MAX]),
1950            DataArrayValues::I16(vec![i16::MIN, i16::MAX]),
1951            DataArrayValues::I32(vec![i32::MIN, i32::MAX]),
1952            DataArrayValues::I64(vec![i64::MIN, i64::MAX]),
1953            DataArrayValues::U8(vec![0, u8::MAX]),
1954            DataArrayValues::U16(vec![0, u16::MAX]),
1955            DataArrayValues::U32(vec![0, u32::MAX]),
1956            DataArrayValues::U64(vec![0, u64::MAX]),
1957        ];
1958
1959        for values in cases {
1960            let dtype = match &values {
1961                DataArrayValues::I8(_) => "int8",
1962                DataArrayValues::I16(_) => "int16",
1963                DataArrayValues::I32(_) => "int32",
1964                DataArrayValues::I64(_) => "int64",
1965                DataArrayValues::U8(_) => "uint8",
1966                DataArrayValues::U16(_) => "uint16",
1967                DataArrayValues::U32(_) => "uint32",
1968                DataArrayValues::U64(_) => "uint64",
1969                DataArrayValues::F64(_) | DataArrayValues::F32(_) => unreachable!(),
1970            };
1971            let payload = DataArrayPayload {
1972                dtype: dtype.to_string(),
1973                shape: vec![1, 2],
1974                values: values.clone(),
1975                imaginary_values: None,
1976            };
1977            let bytes = serde_json::to_vec(&payload).expect("encode typed payload");
1978            let decoded: DataArrayPayload = serde_json::from_slice(&bytes).expect("decode payload");
1979            assert_eq!(decoded.values, values, "{dtype} payload must remain exact");
1980            let Value::Tensor(tensor) = decoded.into_value().expect("tensor value") else {
1981                panic!("expected tensor");
1982            };
1983            assert_eq!(
1984                tensor.integer_storage().map(IntegerStorage::class_name),
1985                Some(dtype)
1986            );
1987        }
1988    }
1989
1990    #[test]
1991    fn payload_roundtrips_native_single_storage() {
1992        let values = DataArrayValues::F32(vec![f32::MIN, 0.1, f32::MAX]);
1993        let payload = DataArrayPayload {
1994            dtype: "f32".to_string(),
1995            shape: vec![1, 3],
1996            values: values.clone(),
1997            imaginary_values: None,
1998        };
1999        let bytes = serde_json::to_vec(&payload).expect("encode single payload");
2000        let decoded: DataArrayPayload =
2001            serde_json::from_slice(&bytes).expect("decode single payload");
2002        assert_eq!(decoded.values, values);
2003
2004        let Value::Tensor(tensor) = decoded.into_value().expect("single tensor value") else {
2005            panic!("expected tensor");
2006        };
2007        assert_eq!(tensor.numeric_dtype(), runmat_value::NumericDType::F32);
2008        assert_eq!(
2009            tensor.materialize_f64(),
2010            vec![f64::from(f32::MIN), f64::from(0.1_f32), f64::from(f32::MAX)]
2011        );
2012    }
2013
2014    #[test]
2015    fn payload_construction_preserves_native_single_tensor() {
2016        let input = Tensor::from_f32(vec![0.1, -2.5], vec![1, 2]).expect("single tensor");
2017        let payload = DataArrayPayload::from_value("f32".to_string(), &Value::Tensor(input))
2018            .expect("single payload");
2019        assert_eq!(payload.values, DataArrayValues::F32(vec![0.1, -2.5]));
2020    }
2021
2022    #[test]
2023    fn payload_decodes_legacy_f64_arrays_and_normalizes_declared_integer_dtypes() {
2024        let legacy = br#"{"dtype":"uint64","shape":[1,2],"values":[1,2]}"#;
2025        let payload: DataArrayPayload =
2026            serde_json::from_slice(legacy).expect("decode legacy payload");
2027        assert_eq!(payload.values, DataArrayValues::F64(vec![1.0, 2.0]));
2028
2029        let payload = payload
2030            .normalize_for_dtype("uint64")
2031            .expect("normalize legacy payload");
2032        assert_eq!(payload.values, DataArrayValues::U64(vec![1, 2]));
2033    }
2034
2035    #[test]
2036    fn preview_conversion_is_bounded_for_typed_integer_payloads() {
2037        let values = DataArrayValues::I16(vec![-2, 0, 3, 7]);
2038
2039        assert_eq!(values.preview_f64(3), vec![-2.0, 0.0, 3.0]);
2040        assert!(values.preview_f64(0).is_empty());
2041    }
2042
2043    #[test]
2044    fn payload_construction_preserves_uint64_tensor_extrema() {
2045        let input =
2046            Tensor::new_integer(IntegerStorage::U64(vec![1_u64 << 63, u64::MAX]), vec![1, 2])
2047                .expect("uint64 tensor");
2048        let payload = DataArrayPayload::from_value("uint64".to_string(), &Value::Tensor(input))
2049            .expect("payload");
2050        assert_eq!(
2051            payload.values,
2052            DataArrayValues::U64(vec![1_u64 << 63, u64::MAX])
2053        );
2054    }
2055
2056    #[test]
2057    fn payload_rejects_typed_complex_integers_without_float_coercion() {
2058        let storage = runmat_value::IntegerComplexStorage::new(
2059            IntegerStorage::U64(vec![1_u64 << 63, u64::MAX]),
2060            IntegerStorage::U64(vec![u64::MAX, 1_u64 << 63]),
2061        )
2062        .expect("matching typed complex components");
2063        let complex = runmat_value::ComplexTensor::new_integer(storage.clone(), vec![1, 2])
2064            .expect("typed complex tensor");
2065
2066        let payload =
2067            DataArrayPayload::from_value("uint64".to_string(), &Value::ComplexTensor(complex))
2068                .expect("encode paired integer payload");
2069        assert_eq!(
2070            payload.imaginary_values,
2071            Some(DataArrayValues::U64(vec![u64::MAX, 1_u64 << 63]))
2072        );
2073        let bytes = serde_json::to_vec(&payload).expect("serialize paired payload");
2074        let decoded: DataArrayPayload =
2075            serde_json::from_slice(&bytes).expect("deserialize paired payload");
2076        let Value::ComplexTensor(decoded) = decoded.into_value().expect("decode paired payload")
2077        else {
2078            panic!("expected paired complex tensor");
2079        };
2080        assert_eq!(decoded.integer_storage(), Some(&storage));
2081    }
2082
2083    #[test]
2084    fn ensure_manifest_sequence_accepts_matching_sequence() {
2085        let manifest = DataManifest {
2086            schema_version: 1,
2087            format: "runmat-data".to_string(),
2088            dataset_id: "ds_test".to_string(),
2089            name: Some("test".to_string()),
2090            created_at: "2026-03-01T00:00:00Z".to_string(),
2091            updated_at: "2026-03-01T00:00:00Z".to_string(),
2092            arrays: BTreeMap::new(),
2093            attrs: BTreeMap::new(),
2094            txn_sequence: 5,
2095        };
2096        ensure_manifest_sequence(5, &manifest).expect("expected sequence match");
2097    }
2098
2099    #[test]
2100    fn ensure_manifest_sequence_rejects_conflict() {
2101        let manifest = DataManifest {
2102            schema_version: 1,
2103            format: "runmat-data".to_string(),
2104            dataset_id: "ds_test".to_string(),
2105            name: Some("test".to_string()),
2106            created_at: "2026-03-01T00:00:00Z".to_string(),
2107            updated_at: "2026-03-01T00:00:00Z".to_string(),
2108            arrays: BTreeMap::new(),
2109            attrs: BTreeMap::new(),
2110            txn_sequence: 6,
2111        };
2112        let err = ensure_manifest_sequence(5, &manifest).expect_err("expected conflict error");
2113        assert_eq!(
2114            err.identifier(),
2115            Some(DATA_MANIFEST_CONFLICT_IDENTIFIER),
2116            "manifest conflicts should expose a stable identifier"
2117        );
2118    }
2119
2120    #[test]
2121    fn transaction_registry_roundtrip() {
2122        let tx_id = start_tx("/datasets/test.data".to_string(), 7).expect("start tx");
2123        let status = with_tx(&tx_id, |tx| Ok(tx.status.clone())).expect("tx lookup");
2124        assert_eq!(status, TxnStatus::Open);
2125        remove_tx(&tx_id).expect("remove tx");
2126        let err = with_tx(&tx_id, |_| Ok(())).expect_err("expected missing tx");
2127        assert_eq!(
2128            err.identifier(),
2129            Some(DATA_TRANSACTION_NOT_FOUND_IDENTIFIER),
2130            "missing transaction lookups should expose a stable identifier"
2131        );
2132    }
2133
2134    #[cfg(not(target_arch = "wasm32"))]
2135    #[tokio::test(flavor = "multi_thread", worker_threads = 2)]
2136    async fn transaction_registry_scope_survives_await() {
2137        with_tx_registry_scope(async {
2138            let tx_id = start_tx("/datasets/task-local.data".to_string(), 11).expect("start tx");
2139            tokio::task::yield_now().await;
2140            let status = with_tx(&tx_id, |tx| Ok(tx.status.clone())).expect("tx lookup");
2141            assert_eq!(status, TxnStatus::Open);
2142            remove_tx(&tx_id).expect("remove tx");
2143            let err = with_tx(&tx_id, |_| Ok(())).expect_err("expected missing tx");
2144            assert_eq!(
2145                err.identifier(),
2146                Some(DATA_TRANSACTION_NOT_FOUND_IDENTIFIER)
2147            );
2148        })
2149        .await;
2150    }
2151
2152    #[test]
2153    fn sha256_hash_format_matches_expected_prefix() {
2154        let hash = sha256_hex(b"runmat");
2155        assert!(hash.starts_with("sha256:"));
2156        assert_eq!(hash.len(), "sha256:".len() + 64);
2157    }
2158}