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#[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 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}