Skip to main content

runmat_runtime/builtins/io/data/
mod.rs

1//! Cloud-ready dataset persistence builtins (`data.*`).
2
3use std::collections::BTreeMap;
4use std::collections::HashMap;
5use std::path::PathBuf;
6
7use runmat_builtins::{
8    BuiltinCompletionPolicy, BuiltinDescriptor, BuiltinErrorDescriptor, BuiltinOutputMode,
9    BuiltinParamArity, BuiltinParamDescriptor, BuiltinParamType, BuiltinSignatureDescriptor,
10    ObjectInstance, StructValue, Tensor, Value,
11};
12use runmat_filesystem::data_contract::{DataChunkDescriptor, DataChunkUploadRequest};
13use runmat_macros::runtime_builtin;
14
15use crate::builtins::common::spec::{
16    BroadcastSemantics, BuiltinFusionSpec, BuiltinGpuSpec, ConstantStrategy, GpuOpKind,
17    ReductionNaN, ResidencyPolicy, ShapeRequirements,
18};
19use crate::data::{
20    array_object, data_error, dataset_object, dataset_root, ensure_manifest_sequence,
21    get_object_prop, manifest_path, manifest_version_token, now_rfc3339, parse_schema,
22    parse_string, read_array_payload_async, read_manifest_async, remove_tx, sha256_hex, start_tx,
23    transaction_object, with_tx, with_tx_mut, write_array_payload_async, write_manifest_async,
24    DataArrayMeta, DataArrayPayload, DataChunkIndex, DataChunkIndexEntry, DataManifest,
25    PendingCreateArray, PendingFill, PendingResize, PendingWrite, TxnStatus,
26};
27use crate::{make_cell, BuiltinResult};
28
29#[runmat_macros::register_gpu_spec(builtin_path = "crate::builtins::io::data")]
30pub const GPU_SPEC: BuiltinGpuSpec = BuiltinGpuSpec {
31    name: "data.*",
32    op_kind: GpuOpKind::Custom("io-data"),
33    supported_precisions: &[],
34    broadcast: BroadcastSemantics::None,
35    provider_hooks: &[],
36    constant_strategy: ConstantStrategy::InlineLiteral,
37    residency: ResidencyPolicy::GatherImmediately,
38    nan_mode: ReductionNaN::Include,
39    two_pass_threshold: None,
40    workgroup_size: None,
41    accepts_nan_mode: false,
42    notes: "Dataset operations are host I/O and metadata orchestration.",
43};
44
45#[runmat_macros::register_fusion_spec(builtin_path = "crate::builtins::io::data")]
46pub const FUSION_SPEC: BuiltinFusionSpec = BuiltinFusionSpec {
47    name: "data.*",
48    shape: ShapeRequirements::Any,
49    constant_strategy: ConstantStrategy::InlineLiteral,
50    elementwise: None,
51    reduction: None,
52    emits_nan: false,
53    notes: "Data builtins are side-effecting and not fusible.",
54};
55
56const DATA_ERROR_INVALID_ARGUMENT: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
57    code: "RM.DATA.INVALID_ARGUMENT",
58    identifier: Some("RunMat:data:InvalidArgument"),
59    when: "Arguments, receiver object, or option grammar are invalid for the requested data API.",
60    message: "data: invalid argument",
61};
62
63const DATA_ERROR_NOT_FOUND: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
64    code: "RM.DATA.NOT_FOUND",
65    identifier: Some("RunMat:data:NotFound"),
66    when: "Referenced dataset/array/transaction object is missing or cannot be resolved.",
67    message: "data: requested object not found",
68};
69
70const DATA_ERROR_INTERNAL: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
71    code: "RM.DATA.INTERNAL",
72    identifier: Some("RunMat:data:Internal"),
73    when: "Filesystem/manifest/chunk processing fails unexpectedly during data API execution.",
74    message: "data: internal operation failed",
75};
76
77const DATA_DESCRIPTOR_ERRORS: [BuiltinErrorDescriptor; 3] = [
78    DATA_ERROR_INVALID_ARGUMENT,
79    DATA_ERROR_NOT_FOUND,
80    DATA_ERROR_INTERNAL,
81];
82
83const OUT_DATASET: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
84    name: "ds",
85    ty: BuiltinParamType::Any,
86    arity: BuiltinParamArity::Required,
87    default: None,
88    description: "Dataset handle object.",
89}];
90
91const OUT_ARRAY: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
92    name: "arr",
93    ty: BuiltinParamType::Any,
94    arity: BuiltinParamArity::Required,
95    default: None,
96    description: "DataArray handle object.",
97}];
98
99const OUT_TX: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
100    name: "tx",
101    ty: BuiltinParamType::Any,
102    arity: BuiltinParamArity::Required,
103    default: None,
104    description: "DataTransaction handle object.",
105}];
106
107const OUT_BOOL: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
108    name: "ok",
109    ty: BuiltinParamType::LogicalArray,
110    arity: BuiltinParamArity::Required,
111    default: None,
112    description: "Logical success/result flag.",
113}];
114
115const OUT_STRING: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
116    name: "s",
117    ty: BuiltinParamType::StringScalar,
118    arity: BuiltinParamArity::Required,
119    default: None,
120    description: "String scalar result.",
121}];
122
123const OUT_STRUCT: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
124    name: "S",
125    ty: BuiltinParamType::Any,
126    arity: BuiltinParamArity::Required,
127    default: None,
128    description: "Struct result.",
129}];
130
131const OUT_CELL: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
132    name: "C",
133    ty: BuiltinParamType::Any,
134    arity: BuiltinParamArity::Required,
135    default: None,
136    description: "Cell-array result.",
137}];
138
139const OUT_VALUE: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
140    name: "value",
141    ty: BuiltinParamType::Any,
142    arity: BuiltinParamArity::Required,
143    default: None,
144    description: "Value result.",
145}];
146
147const OUT_TENSOR: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
148    name: "X",
149    ty: BuiltinParamType::NumericArray,
150    arity: BuiltinParamArity::Required,
151    default: None,
152    description: "Numeric tensor result.",
153}];
154
155const IN_PATH_SCHEMA_REST: [BuiltinParamDescriptor; 3] = [
156    BuiltinParamDescriptor {
157        name: "path",
158        ty: BuiltinParamType::StringScalar,
159        arity: BuiltinParamArity::Required,
160        default: None,
161        description: "Dataset path (.data).",
162    },
163    BuiltinParamDescriptor {
164        name: "schema",
165        ty: BuiltinParamType::Any,
166        arity: BuiltinParamArity::Required,
167        default: None,
168        description: "Dataset schema struct.",
169    },
170    BuiltinParamDescriptor {
171        name: "options",
172        ty: BuiltinParamType::Any,
173        arity: BuiltinParamArity::Variadic,
174        default: None,
175        description: "Reserved name/value options.",
176    },
177];
178
179const IN_PATH_REST: [BuiltinParamDescriptor; 2] = [
180    BuiltinParamDescriptor {
181        name: "path",
182        ty: BuiltinParamType::StringScalar,
183        arity: BuiltinParamArity::Required,
184        default: None,
185        description: "Dataset path (.data).",
186    },
187    BuiltinParamDescriptor {
188        name: "options",
189        ty: BuiltinParamType::Any,
190        arity: BuiltinParamArity::Variadic,
191        default: None,
192        description: "Reserved name/value options.",
193    },
194];
195
196const IN_FROM_TO_REST: [BuiltinParamDescriptor; 3] = [
197    BuiltinParamDescriptor {
198        name: "fromPath",
199        ty: BuiltinParamType::StringScalar,
200        arity: BuiltinParamArity::Required,
201        default: None,
202        description: "Source dataset path.",
203    },
204    BuiltinParamDescriptor {
205        name: "toPath",
206        ty: BuiltinParamType::StringScalar,
207        arity: BuiltinParamArity::Required,
208        default: None,
209        description: "Destination dataset path.",
210    },
211    BuiltinParamDescriptor {
212        name: "options",
213        ty: BuiltinParamType::Any,
214        arity: BuiltinParamArity::Variadic,
215        default: None,
216        description: "Reserved name/value options.",
217    },
218];
219
220const IN_PATH_FORMAT_TARGET_REST: [BuiltinParamDescriptor; 4] = [
221    BuiltinParamDescriptor {
222        name: "path",
223        ty: BuiltinParamType::StringScalar,
224        arity: BuiltinParamArity::Required,
225        default: None,
226        description: "Dataset path (.data).",
227    },
228    BuiltinParamDescriptor {
229        name: "format",
230        ty: BuiltinParamType::StringScalar,
231        arity: BuiltinParamArity::Required,
232        default: None,
233        description: "Format token (currently 'data').",
234    },
235    BuiltinParamDescriptor {
236        name: "targetPath",
237        ty: BuiltinParamType::StringScalar,
238        arity: BuiltinParamArity::Required,
239        default: None,
240        description: "Target dataset path.",
241    },
242    BuiltinParamDescriptor {
243        name: "options",
244        ty: BuiltinParamType::Any,
245        arity: BuiltinParamArity::Variadic,
246        default: None,
247        description: "Reserved name/value options.",
248    },
249];
250
251const IN_PREFIX_REST: [BuiltinParamDescriptor; 2] = [
252    BuiltinParamDescriptor {
253        name: "prefix",
254        ty: BuiltinParamType::StringScalar,
255        arity: BuiltinParamArity::Required,
256        default: None,
257        description: "Filesystem prefix to scan for datasets.",
258    },
259    BuiltinParamDescriptor {
260        name: "options",
261        ty: BuiltinParamType::Any,
262        arity: BuiltinParamArity::Variadic,
263        default: None,
264        description: "Reserved name/value options.",
265    },
266];
267
268const IN_BASE: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
269    name: "obj",
270    ty: BuiltinParamType::Any,
271    arity: BuiltinParamArity::Required,
272    default: None,
273    description: "Receiver object.",
274}];
275
276const IN_BASE_NAME: [BuiltinParamDescriptor; 2] = [
277    BuiltinParamDescriptor {
278        name: "obj",
279        ty: BuiltinParamType::Any,
280        arity: BuiltinParamArity::Required,
281        default: None,
282        description: "Receiver object.",
283    },
284    BuiltinParamDescriptor {
285        name: "name",
286        ty: BuiltinParamType::StringScalar,
287        arity: BuiltinParamArity::Required,
288        default: None,
289        description: "Name key/array identifier.",
290    },
291];
292
293const IN_BASE_KEY_DEFAULT: [BuiltinParamDescriptor; 3] = [
294    BuiltinParamDescriptor {
295        name: "obj",
296        ty: BuiltinParamType::Any,
297        arity: BuiltinParamArity::Required,
298        default: None,
299        description: "Receiver object.",
300    },
301    BuiltinParamDescriptor {
302        name: "key",
303        ty: BuiltinParamType::StringScalar,
304        arity: BuiltinParamArity::Required,
305        default: None,
306        description: "Attribute key.",
307    },
308    BuiltinParamDescriptor {
309        name: "defaultValue",
310        ty: BuiltinParamType::Any,
311        arity: BuiltinParamArity::Optional,
312        default: Some("0"),
313        description: "Fallback value when key is absent.",
314    },
315];
316
317const IN_BASE_KEY_VALUE: [BuiltinParamDescriptor; 3] = [
318    BuiltinParamDescriptor {
319        name: "obj",
320        ty: BuiltinParamType::Any,
321        arity: BuiltinParamArity::Required,
322        default: None,
323        description: "Receiver object.",
324    },
325    BuiltinParamDescriptor {
326        name: "key",
327        ty: BuiltinParamType::StringScalar,
328        arity: BuiltinParamArity::Required,
329        default: None,
330        description: "Attribute key.",
331    },
332    BuiltinParamDescriptor {
333        name: "value",
334        ty: BuiltinParamType::Any,
335        arity: BuiltinParamArity::Required,
336        default: None,
337        description: "Attribute value.",
338    },
339];
340
341const IN_BASE_ATTRS: [BuiltinParamDescriptor; 2] = [
342    BuiltinParamDescriptor {
343        name: "obj",
344        ty: BuiltinParamType::Any,
345        arity: BuiltinParamArity::Required,
346        default: None,
347        description: "Receiver object.",
348    },
349    BuiltinParamDescriptor {
350        name: "attrs",
351        ty: BuiltinParamType::Any,
352        arity: BuiltinParamArity::Required,
353        default: None,
354        description: "Struct of attribute updates.",
355    },
356];
357
358const IN_BASE_LABEL_REST: [BuiltinParamDescriptor; 3] = [
359    BuiltinParamDescriptor {
360        name: "obj",
361        ty: BuiltinParamType::Any,
362        arity: BuiltinParamArity::Required,
363        default: None,
364        description: "Receiver object.",
365    },
366    BuiltinParamDescriptor {
367        name: "label",
368        ty: BuiltinParamType::StringScalar,
369        arity: BuiltinParamArity::Required,
370        default: None,
371        description: "Snapshot label.",
372    },
373    BuiltinParamDescriptor {
374        name: "options",
375        ty: BuiltinParamType::Any,
376        arity: BuiltinParamArity::Variadic,
377        default: None,
378        description: "Reserved name/value options.",
379    },
380];
381
382const IN_BASE_REST: [BuiltinParamDescriptor; 2] = [
383    BuiltinParamDescriptor {
384        name: "obj",
385        ty: BuiltinParamType::Any,
386        arity: BuiltinParamArity::Required,
387        default: None,
388        description: "Receiver object.",
389    },
390    BuiltinParamDescriptor {
391        name: "options",
392        ty: BuiltinParamType::Any,
393        arity: BuiltinParamArity::Variadic,
394        default: None,
395        description: "Reserved name/value options.",
396    },
397];
398
399const IN_BASE_SLICE_OPTIONAL: [BuiltinParamDescriptor; 2] = [
400    BuiltinParamDescriptor {
401        name: "obj",
402        ty: BuiltinParamType::Any,
403        arity: BuiltinParamArity::Required,
404        default: None,
405        description: "DataArray receiver object.",
406    },
407    BuiltinParamDescriptor {
408        name: "sliceSpec",
409        ty: BuiltinParamType::Any,
410        arity: BuiltinParamArity::Optional,
411        default: None,
412        description: "Optional slice specification.",
413    },
414];
415
416const IN_BASE_VALUES: [BuiltinParamDescriptor; 2] = [
417    BuiltinParamDescriptor {
418        name: "obj",
419        ty: BuiltinParamType::Any,
420        arity: BuiltinParamArity::Required,
421        default: None,
422        description: "DataArray receiver object.",
423    },
424    BuiltinParamDescriptor {
425        name: "values",
426        ty: BuiltinParamType::Any,
427        arity: BuiltinParamArity::Required,
428        default: None,
429        description: "Full-array values payload.",
430    },
431];
432
433const IN_BASE_SLICE_VALUES: [BuiltinParamDescriptor; 3] = [
434    BuiltinParamDescriptor {
435        name: "obj",
436        ty: BuiltinParamType::Any,
437        arity: BuiltinParamArity::Required,
438        default: None,
439        description: "DataArray receiver object.",
440    },
441    BuiltinParamDescriptor {
442        name: "sliceSpec",
443        ty: BuiltinParamType::Any,
444        arity: BuiltinParamArity::Required,
445        default: None,
446        description: "Slice specification.",
447    },
448    BuiltinParamDescriptor {
449        name: "values",
450        ty: BuiltinParamType::Any,
451        arity: BuiltinParamArity::Required,
452        default: None,
453        description: "Slice values payload.",
454    },
455];
456
457const IN_BASE_NEW_SHAPE_REST: [BuiltinParamDescriptor; 3] = [
458    BuiltinParamDescriptor {
459        name: "obj",
460        ty: BuiltinParamType::Any,
461        arity: BuiltinParamArity::Required,
462        default: None,
463        description: "DataArray receiver object.",
464    },
465    BuiltinParamDescriptor {
466        name: "newShape",
467        ty: BuiltinParamType::Any,
468        arity: BuiltinParamArity::Required,
469        default: None,
470        description: "New array shape.",
471    },
472    BuiltinParamDescriptor {
473        name: "options",
474        ty: BuiltinParamType::Any,
475        arity: BuiltinParamArity::Variadic,
476        default: None,
477        description: "Reserved name/value options.",
478    },
479];
480
481const IN_BASE_VALUE_REST: [BuiltinParamDescriptor; 3] = [
482    BuiltinParamDescriptor {
483        name: "obj",
484        ty: BuiltinParamType::Any,
485        arity: BuiltinParamArity::Required,
486        default: None,
487        description: "DataArray receiver object.",
488    },
489    BuiltinParamDescriptor {
490        name: "value",
491        ty: BuiltinParamType::Any,
492        arity: BuiltinParamArity::Required,
493        default: None,
494        description: "Fill value.",
495    },
496    BuiltinParamDescriptor {
497        name: "options",
498        ty: BuiltinParamType::Any,
499        arity: BuiltinParamArity::Variadic,
500        default: None,
501        description: "Reserved name/value options.",
502    },
503];
504
505const IN_TX_WRITE: [BuiltinParamDescriptor; 5] = [
506    BuiltinParamDescriptor {
507        name: "tx",
508        ty: BuiltinParamType::Any,
509        arity: BuiltinParamArity::Required,
510        default: None,
511        description: "DataTransaction receiver.",
512    },
513    BuiltinParamDescriptor {
514        name: "arrayName",
515        ty: BuiltinParamType::StringScalar,
516        arity: BuiltinParamArity::Required,
517        default: None,
518        description: "Target array name.",
519    },
520    BuiltinParamDescriptor {
521        name: "sliceSpec",
522        ty: BuiltinParamType::Any,
523        arity: BuiltinParamArity::Required,
524        default: None,
525        description: "Slice specification.",
526    },
527    BuiltinParamDescriptor {
528        name: "values",
529        ty: BuiltinParamType::Any,
530        arity: BuiltinParamArity::Required,
531        default: None,
532        description: "Values payload.",
533    },
534    BuiltinParamDescriptor {
535        name: "options",
536        ty: BuiltinParamType::Any,
537        arity: BuiltinParamArity::Variadic,
538        default: None,
539        description: "Reserved name/value options.",
540    },
541];
542
543const IN_TX_ARRAY_SHAPE_REST: [BuiltinParamDescriptor; 4] = [
544    BuiltinParamDescriptor {
545        name: "tx",
546        ty: BuiltinParamType::Any,
547        arity: BuiltinParamArity::Required,
548        default: None,
549        description: "DataTransaction receiver.",
550    },
551    BuiltinParamDescriptor {
552        name: "arrayName",
553        ty: BuiltinParamType::StringScalar,
554        arity: BuiltinParamArity::Required,
555        default: None,
556        description: "Target array name.",
557    },
558    BuiltinParamDescriptor {
559        name: "newShape",
560        ty: BuiltinParamType::Any,
561        arity: BuiltinParamArity::Required,
562        default: None,
563        description: "New shape vector.",
564    },
565    BuiltinParamDescriptor {
566        name: "options",
567        ty: BuiltinParamType::Any,
568        arity: BuiltinParamArity::Variadic,
569        default: None,
570        description: "Reserved name/value options.",
571    },
572];
573
574const IN_TX_ARRAY_VALUE_SLICE_OPT: [BuiltinParamDescriptor; 4] = [
575    BuiltinParamDescriptor {
576        name: "tx",
577        ty: BuiltinParamType::Any,
578        arity: BuiltinParamArity::Required,
579        default: None,
580        description: "DataTransaction receiver.",
581    },
582    BuiltinParamDescriptor {
583        name: "arrayName",
584        ty: BuiltinParamType::StringScalar,
585        arity: BuiltinParamArity::Required,
586        default: None,
587        description: "Target array name.",
588    },
589    BuiltinParamDescriptor {
590        name: "value",
591        ty: BuiltinParamType::Any,
592        arity: BuiltinParamArity::Required,
593        default: None,
594        description: "Fill value.",
595    },
596    BuiltinParamDescriptor {
597        name: "sliceSpec",
598        ty: BuiltinParamType::Any,
599        arity: BuiltinParamArity::Optional,
600        default: None,
601        description: "Optional slice specification.",
602    },
603];
604
605const IN_TX_ARRAY_NAME: [BuiltinParamDescriptor; 2] = [
606    BuiltinParamDescriptor {
607        name: "tx",
608        ty: BuiltinParamType::Any,
609        arity: BuiltinParamArity::Required,
610        default: None,
611        description: "DataTransaction receiver.",
612    },
613    BuiltinParamDescriptor {
614        name: "arrayName",
615        ty: BuiltinParamType::StringScalar,
616        arity: BuiltinParamArity::Required,
617        default: None,
618        description: "Target array name.",
619    },
620];
621
622const IN_TX_ARRAY_META: [BuiltinParamDescriptor; 3] = [
623    BuiltinParamDescriptor {
624        name: "tx",
625        ty: BuiltinParamType::Any,
626        arity: BuiltinParamArity::Required,
627        default: None,
628        description: "DataTransaction receiver.",
629    },
630    BuiltinParamDescriptor {
631        name: "arrayName",
632        ty: BuiltinParamType::StringScalar,
633        arity: BuiltinParamArity::Required,
634        default: None,
635        description: "New array name.",
636    },
637    BuiltinParamDescriptor {
638        name: "meta",
639        ty: BuiltinParamType::Any,
640        arity: BuiltinParamArity::Required,
641        default: None,
642        description: "Array metadata struct.",
643    },
644];
645
646const IN_TX_COMMIT_REST: [BuiltinParamDescriptor; 2] = [
647    BuiltinParamDescriptor {
648        name: "tx",
649        ty: BuiltinParamType::Any,
650        arity: BuiltinParamArity::Required,
651        default: None,
652        description: "DataTransaction receiver.",
653    },
654    BuiltinParamDescriptor {
655        name: "options",
656        ty: BuiltinParamType::Any,
657        arity: BuiltinParamArity::Variadic,
658        default: None,
659        description: "Commit options (e.g. if_manifest).",
660    },
661];
662
663macro_rules! one_sig_descriptor {
664    ($desc:ident, $sigs:ident, $label:expr, $inputs:expr, $outputs:expr) => {
665        const $sigs: [BuiltinSignatureDescriptor; 1] = [BuiltinSignatureDescriptor {
666            label: $label,
667            inputs: $inputs,
668            outputs: $outputs,
669        }];
670        pub const $desc: BuiltinDescriptor = BuiltinDescriptor {
671            signatures: &$sigs,
672            output_mode: BuiltinOutputMode::Fixed,
673            completion_policy: BuiltinCompletionPolicy::Public,
674            errors: &DATA_DESCRIPTOR_ERRORS,
675        };
676    };
677}
678
679one_sig_descriptor!(
680    DATA_CREATE_DESCRIPTOR,
681    DATA_CREATE_SIGS,
682    "ds = data.create(path, schema, Name, Value, ...)",
683    &IN_PATH_SCHEMA_REST,
684    &OUT_DATASET
685);
686one_sig_descriptor!(
687    DATA_OPEN_DESCRIPTOR,
688    DATA_OPEN_SIGS,
689    "ds = data.open(path, Name, Value, ...)",
690    &IN_PATH_REST,
691    &OUT_DATASET
692);
693one_sig_descriptor!(
694    DATA_EXISTS_DESCRIPTOR,
695    DATA_EXISTS_SIGS,
696    "tf = data.exists(path)",
697    &IN_PATH_REST,
698    &OUT_BOOL
699);
700one_sig_descriptor!(
701    DATA_DELETE_DESCRIPTOR,
702    DATA_DELETE_SIGS,
703    "tf = data.delete(path, Name, Value, ...)",
704    &IN_PATH_REST,
705    &OUT_BOOL
706);
707one_sig_descriptor!(
708    DATA_COPY_DESCRIPTOR,
709    DATA_COPY_SIGS,
710    "tf = data.copy(fromPath, toPath, Name, Value, ...)",
711    &IN_FROM_TO_REST,
712    &OUT_BOOL
713);
714one_sig_descriptor!(
715    DATA_MOVE_DESCRIPTOR,
716    DATA_MOVE_SIGS,
717    "tf = data.move(fromPath, toPath, Name, Value, ...)",
718    &IN_FROM_TO_REST,
719    &OUT_BOOL
720);
721one_sig_descriptor!(
722    DATA_IMPORT_DESCRIPTOR,
723    DATA_IMPORT_SIGS,
724    "ds = data.import(path, format, sourcePath, Name, Value, ...)",
725    &IN_PATH_FORMAT_TARGET_REST,
726    &OUT_DATASET
727);
728one_sig_descriptor!(
729    DATA_EXPORT_DESCRIPTOR,
730    DATA_EXPORT_SIGS,
731    "tf = data.export(path, format, targetPath, Name, Value, ...)",
732    &IN_PATH_FORMAT_TARGET_REST,
733    &OUT_BOOL
734);
735one_sig_descriptor!(
736    DATA_LIST_DESCRIPTOR,
737    DATA_LIST_SIGS,
738    "C = data.list(prefix, Name, Value, ...)",
739    &IN_PREFIX_REST,
740    &OUT_CELL
741);
742one_sig_descriptor!(
743    DATA_INSPECT_DESCRIPTOR,
744    DATA_INSPECT_SIGS,
745    "S = data.inspect(path)",
746    &IN_PATH_REST,
747    &OUT_STRUCT
748);
749
750one_sig_descriptor!(
751    DATASET_PATH_DESCRIPTOR,
752    DATASET_PATH_SIGS,
753    "path = Dataset.path(ds)",
754    &IN_BASE,
755    &OUT_STRING
756);
757one_sig_descriptor!(
758    DATASET_ID_DESCRIPTOR,
759    DATASET_ID_SIGS,
760    "id = Dataset.id(ds)",
761    &IN_BASE,
762    &OUT_STRING
763);
764one_sig_descriptor!(
765    DATASET_VERSION_DESCRIPTOR,
766    DATASET_VERSION_SIGS,
767    "version = Dataset.version(ds)",
768    &IN_BASE,
769    &OUT_STRING
770);
771one_sig_descriptor!(
772    DATASET_ARRAYS_DESCRIPTOR,
773    DATASET_ARRAYS_SIGS,
774    "C = Dataset.arrays(ds)",
775    &IN_BASE,
776    &OUT_CELL
777);
778one_sig_descriptor!(
779    DATASET_HAS_ARRAY_DESCRIPTOR,
780    DATASET_HAS_ARRAY_SIGS,
781    "tf = Dataset.has_array(ds, name)",
782    &IN_BASE_NAME,
783    &OUT_BOOL
784);
785one_sig_descriptor!(
786    DATASET_ARRAY_DESCRIPTOR,
787    DATASET_ARRAY_SIGS,
788    "arr = Dataset.array(ds, name)",
789    &IN_BASE_NAME,
790    &OUT_ARRAY
791);
792one_sig_descriptor!(
793    DATASET_ATTRS_DESCRIPTOR,
794    DATASET_ATTRS_SIGS,
795    "S = Dataset.attrs(ds)",
796    &IN_BASE,
797    &OUT_STRUCT
798);
799one_sig_descriptor!(
800    DATASET_GET_ATTR_DESCRIPTOR,
801    DATASET_GET_ATTR_SIGS,
802    "value = Dataset.get_attr(ds, key, defaultValue)",
803    &IN_BASE_KEY_DEFAULT,
804    &OUT_VALUE
805);
806one_sig_descriptor!(
807    DATASET_SET_ATTR_DESCRIPTOR,
808    DATASET_SET_ATTR_SIGS,
809    "tf = Dataset.set_attr(ds, key, value)",
810    &IN_BASE_KEY_VALUE,
811    &OUT_BOOL
812);
813one_sig_descriptor!(
814    DATASET_SET_ATTRS_DESCRIPTOR,
815    DATASET_SET_ATTRS_SIGS,
816    "tf = Dataset.set_attrs(ds, attrs)",
817    &IN_BASE_ATTRS,
818    &OUT_BOOL
819);
820one_sig_descriptor!(
821    DATASET_BEGIN_DESCRIPTOR,
822    DATASET_BEGIN_SIGS,
823    "tx = Dataset.begin(ds, Name, Value, ...)",
824    &IN_BASE_REST,
825    &OUT_TX
826);
827one_sig_descriptor!(
828    DATASET_SNAPSHOT_DESCRIPTOR,
829    DATASET_SNAPSHOT_SIGS,
830    "snapshotPath = Dataset.snapshot(ds, label, Name, Value, ...)",
831    &IN_BASE_LABEL_REST,
832    &OUT_STRING
833);
834one_sig_descriptor!(
835    DATASET_REFRESH_DESCRIPTOR,
836    DATASET_REFRESH_SIGS,
837    "ds = Dataset.refresh(ds)",
838    &IN_BASE,
839    &OUT_DATASET
840);
841
842one_sig_descriptor!(
843    DATAARRAY_NAME_DESCRIPTOR,
844    DATAARRAY_NAME_SIGS,
845    "name = DataArray.name(arr)",
846    &IN_BASE,
847    &OUT_STRING
848);
849one_sig_descriptor!(
850    DATAARRAY_DTYPE_DESCRIPTOR,
851    DATAARRAY_DTYPE_SIGS,
852    "dtype = DataArray.dtype(arr)",
853    &IN_BASE,
854    &OUT_STRING
855);
856one_sig_descriptor!(
857    DATAARRAY_SHAPE_DESCRIPTOR,
858    DATAARRAY_SHAPE_SIGS,
859    "shape = DataArray.shape(arr)",
860    &IN_BASE,
861    &OUT_TENSOR
862);
863one_sig_descriptor!(
864    DATAARRAY_RANK_DESCRIPTOR,
865    DATAARRAY_RANK_SIGS,
866    "rank = DataArray.rank(arr)",
867    &IN_BASE,
868    &OUT_VALUE
869);
870one_sig_descriptor!(
871    DATAARRAY_CHUNK_SHAPE_DESCRIPTOR,
872    DATAARRAY_CHUNK_SHAPE_SIGS,
873    "chunkShape = DataArray.chunk_shape(arr)",
874    &IN_BASE,
875    &OUT_TENSOR
876);
877one_sig_descriptor!(
878    DATAARRAY_CODEC_DESCRIPTOR,
879    DATAARRAY_CODEC_SIGS,
880    "codec = DataArray.codec(arr)",
881    &IN_BASE,
882    &OUT_STRING
883);
884one_sig_descriptor!(
885    DATAARRAY_READ_DESCRIPTOR,
886    DATAARRAY_READ_SIGS,
887    "X = DataArray.read(arr, sliceSpec)",
888    &IN_BASE_SLICE_OPTIONAL,
889    &OUT_TENSOR
890);
891const DATAARRAY_WRITE_SIGS: [BuiltinSignatureDescriptor; 2] = [
892    BuiltinSignatureDescriptor {
893        label: "tf = DataArray.write(arr, values)",
894        inputs: &IN_BASE_VALUES,
895        outputs: &OUT_BOOL,
896    },
897    BuiltinSignatureDescriptor {
898        label: "tf = DataArray.write(arr, sliceSpec, values)",
899        inputs: &IN_BASE_SLICE_VALUES,
900        outputs: &OUT_BOOL,
901    },
902];
903pub const DATAARRAY_WRITE_DESCRIPTOR: BuiltinDescriptor = BuiltinDescriptor {
904    signatures: &DATAARRAY_WRITE_SIGS,
905    output_mode: BuiltinOutputMode::Fixed,
906    completion_policy: BuiltinCompletionPolicy::Public,
907    errors: &DATA_DESCRIPTOR_ERRORS,
908};
909one_sig_descriptor!(
910    DATAARRAY_RESIZE_DESCRIPTOR,
911    DATAARRAY_RESIZE_SIGS,
912    "tf = DataArray.resize(arr, newShape, Name, Value, ...)",
913    &IN_BASE_NEW_SHAPE_REST,
914    &OUT_BOOL
915);
916one_sig_descriptor!(
917    DATAARRAY_FILL_DESCRIPTOR,
918    DATAARRAY_FILL_SIGS,
919    "tf = DataArray.fill(arr, value, Name, Value, ...)",
920    &IN_BASE_VALUE_REST,
921    &OUT_BOOL
922);
923
924one_sig_descriptor!(
925    DATATX_ID_DESCRIPTOR,
926    DATATX_ID_SIGS,
927    "id = DataTransaction.id(tx)",
928    &IN_BASE,
929    &OUT_STRING
930);
931one_sig_descriptor!(
932    DATATX_WRITE_DESCRIPTOR,
933    DATATX_WRITE_SIGS_1,
934    "tf = DataTransaction.write(tx, arrayName, sliceSpec, values, Name, Value, ...)",
935    &IN_TX_WRITE,
936    &OUT_BOOL
937);
938one_sig_descriptor!(
939    DATATX_SET_ATTR_DESCRIPTOR,
940    DATATX_SET_ATTR_SIGS,
941    "tf = DataTransaction.set_attr(tx, key, value)",
942    &IN_BASE_KEY_VALUE,
943    &OUT_BOOL
944);
945one_sig_descriptor!(
946    DATATX_SET_ATTRS_DESCRIPTOR,
947    DATATX_SET_ATTRS_SIGS,
948    "tf = DataTransaction.set_attrs(tx, attrs)",
949    &IN_BASE_ATTRS,
950    &OUT_BOOL
951);
952one_sig_descriptor!(
953    DATATX_RESIZE_DESCRIPTOR,
954    DATATX_RESIZE_SIGS,
955    "tf = DataTransaction.resize(tx, arrayName, newShape, Name, Value, ...)",
956    &IN_TX_ARRAY_SHAPE_REST,
957    &OUT_BOOL
958);
959one_sig_descriptor!(
960    DATATX_FILL_DESCRIPTOR,
961    DATATX_FILL_SIGS,
962    "tf = DataTransaction.fill(tx, arrayName, value, sliceSpec)",
963    &IN_TX_ARRAY_VALUE_SLICE_OPT,
964    &OUT_BOOL
965);
966one_sig_descriptor!(
967    DATATX_DELETE_ARRAY_DESCRIPTOR,
968    DATATX_DELETE_ARRAY_SIGS,
969    "tf = DataTransaction.delete_array(tx, arrayName)",
970    &IN_TX_ARRAY_NAME,
971    &OUT_BOOL
972);
973one_sig_descriptor!(
974    DATATX_CREATE_ARRAY_DESCRIPTOR,
975    DATATX_CREATE_ARRAY_SIGS,
976    "tf = DataTransaction.create_array(tx, arrayName, meta)",
977    &IN_TX_ARRAY_META,
978    &OUT_BOOL
979);
980one_sig_descriptor!(
981    DATATX_COMMIT_DESCRIPTOR,
982    DATATX_COMMIT_SIGS,
983    "tf = DataTransaction.commit(tx, Name, Value, ...)",
984    &IN_TX_COMMIT_REST,
985    &OUT_BOOL
986);
987one_sig_descriptor!(
988    COMMIT_ALIAS_DESCRIPTOR,
989    COMMIT_ALIAS_SIGS,
990    "tf = commit(tx, Name, Value, ...)",
991    &IN_TX_COMMIT_REST,
992    &OUT_BOOL
993);
994one_sig_descriptor!(
995    DATATX_ABORT_DESCRIPTOR,
996    DATATX_ABORT_SIGS,
997    "tf = DataTransaction.abort(tx)",
998    &IN_BASE,
999    &OUT_BOOL
1000);
1001one_sig_descriptor!(
1002    DATATX_STATUS_DESCRIPTOR,
1003    DATATX_STATUS_SIGS,
1004    "status = DataTransaction.status(tx)",
1005    &IN_BASE,
1006    &OUT_STRING
1007);
1008
1009#[runtime_builtin(
1010    name = "data.create",
1011    category = "io/data",
1012    summary = "Create a typed dataset at a .data path.",
1013    keywords = "data,dataset,create,persistence",
1014    sink = true,
1015    type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1016    descriptor(crate::builtins::io::data::DATA_CREATE_DESCRIPTOR),
1017    builtin_path = "crate::builtins::io::data"
1018)]
1019async fn data_create_builtin(
1020    path: Value,
1021    schema: Value,
1022    _rest: Vec<Value>,
1023) -> BuiltinResult<Value> {
1024    let path = parse_string(&path, "data.create path")?;
1025    let root = dataset_root(&path);
1026    let schema = parse_schema(&schema)?;
1027    let now = now_rfc3339();
1028    let mut arrays = BTreeMap::new();
1029    for (name, mut meta) in schema.arrays {
1030        let payload = DataArrayPayload::zeros(meta.dtype.clone(), meta.shape.clone());
1031        let (payload_path, chunk_index_path) =
1032            write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
1033        meta.data_path = make_rel_data_path(&root, &payload_path)?;
1034        meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
1035        arrays.insert(name, meta);
1036    }
1037
1038    let manifest = DataManifest {
1039        schema_version: 1,
1040        format: "runmat-data".to_string(),
1041        dataset_id: crate::data::new_dataset_id(),
1042        name: root.file_name().map(|v| v.to_string_lossy().to_string()),
1043        created_at: now.clone(),
1044        updated_at: now,
1045        arrays,
1046        attrs: BTreeMap::new(),
1047        txn_sequence: 0,
1048    };
1049    write_manifest_async(&root, &manifest).await?;
1050    Ok(dataset_object(&path, &manifest))
1051}
1052
1053#[runtime_builtin(
1054    name = "data.open",
1055    category = "io/data",
1056    summary = "Open a dataset handle from a .data path.",
1057    keywords = "data,dataset,open,persistence",
1058    type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1059    descriptor(crate::builtins::io::data::DATA_OPEN_DESCRIPTOR),
1060    builtin_path = "crate::builtins::io::data"
1061)]
1062async fn data_open_builtin(path: Value, _rest: Vec<Value>) -> BuiltinResult<Value> {
1063    let path = parse_string(&path, "data.open path")?;
1064    let root = dataset_root(&path);
1065    let manifest = read_manifest_async(&root).await?;
1066    let mut ds = dataset_object(&path, &manifest);
1067    hydrate_dataset_descriptor_async(&path, &mut ds).await;
1068    Ok(ds)
1069}
1070
1071#[runtime_builtin(
1072    name = "data.exists",
1073    category = "io/data",
1074    summary = "Check if dataset exists.",
1075    keywords = "data,dataset,exists",
1076    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1077    descriptor(crate::builtins::io::data::DATA_EXISTS_DESCRIPTOR),
1078    builtin_path = "crate::builtins::io::data"
1079)]
1080async fn data_exists_builtin(path: Value) -> BuiltinResult<Value> {
1081    let path = parse_string(&path, "data.exists path")?;
1082    let root = dataset_root(&path);
1083    let exists = runmat_filesystem::metadata_async(manifest_path(&root))
1084        .await
1085        .is_ok();
1086    Ok(Value::Bool(exists))
1087}
1088
1089#[runtime_builtin(
1090    name = "data.delete",
1091    category = "io/data",
1092    summary = "Delete a dataset path.",
1093    keywords = "data,dataset,delete",
1094    sink = true,
1095    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1096    descriptor(crate::builtins::io::data::DATA_DELETE_DESCRIPTOR),
1097    builtin_path = "crate::builtins::io::data"
1098)]
1099async fn data_delete_builtin(path: Value, _rest: Vec<Value>) -> BuiltinResult<Value> {
1100    let path = parse_string(&path, "data.delete path")?;
1101    let root = dataset_root(&path);
1102    runmat_filesystem::remove_dir_all_async(&root)
1103        .await
1104        .map_err(|err| {
1105            data_error(format!(
1106                "data.delete: failed to remove '{}': {err}",
1107                root.display()
1108            ))
1109        })?;
1110    Ok(Value::Bool(true))
1111}
1112
1113#[runtime_builtin(
1114    name = "data.copy",
1115    category = "io/data",
1116    summary = "Copy dataset to new path.",
1117    keywords = "data,dataset,copy",
1118    sink = true,
1119    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1120    descriptor(crate::builtins::io::data::DATA_COPY_DESCRIPTOR),
1121    builtin_path = "crate::builtins::io::data"
1122)]
1123async fn data_copy_builtin(
1124    from_path: Value,
1125    to_path: Value,
1126    _rest: Vec<Value>,
1127) -> BuiltinResult<Value> {
1128    let from = parse_string(&from_path, "data.copy fromPath")?;
1129    let to = parse_string(&to_path, "data.copy toPath")?;
1130    copy_dir_recursive(&dataset_root(&from), &dataset_root(&to)).await?;
1131    Ok(Value::Bool(true))
1132}
1133
1134#[runtime_builtin(
1135    name = "data.move",
1136    category = "io/data",
1137    summary = "Move dataset to new path.",
1138    keywords = "data,dataset,move",
1139    sink = true,
1140    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1141    descriptor(crate::builtins::io::data::DATA_MOVE_DESCRIPTOR),
1142    builtin_path = "crate::builtins::io::data"
1143)]
1144async fn data_move_builtin(
1145    from_path: Value,
1146    to_path: Value,
1147    _rest: Vec<Value>,
1148) -> BuiltinResult<Value> {
1149    let from = parse_string(&from_path, "data.move fromPath")?;
1150    let to = parse_string(&to_path, "data.move toPath")?;
1151    runmat_filesystem::rename_async(dataset_root(&from), dataset_root(&to))
1152        .await
1153        .map_err(|err| {
1154            data_error(format!(
1155                "data.move: failed to move dataset '{from}' -> '{to}': {err}"
1156            ))
1157        })?;
1158    Ok(Value::Bool(true))
1159}
1160
1161#[runtime_builtin(
1162    name = "data.import",
1163    category = "io/data",
1164    summary = "Import an existing dataset file path.",
1165    keywords = "data,dataset,import",
1166    sink = true,
1167    type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1168    descriptor(crate::builtins::io::data::DATA_IMPORT_DESCRIPTOR),
1169    builtin_path = "crate::builtins::io::data"
1170)]
1171async fn data_import_builtin(
1172    path: Value,
1173    format: Value,
1174    source_path: Value,
1175    _rest: Vec<Value>,
1176) -> BuiltinResult<Value> {
1177    let path = parse_string(&path, "data.import path")?;
1178    let format = parse_string(&format, "data.import format")?;
1179    if !format.eq_ignore_ascii_case("data") {
1180        return Err(data_error(
1181            "data.import currently supports only format='data'",
1182        ));
1183    }
1184    let source_path = parse_string(&source_path, "data.import sourcePath")?;
1185    copy_dir_recursive(&dataset_root(&source_path), &dataset_root(&path)).await?;
1186    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1187    Ok(dataset_object(&path, &manifest))
1188}
1189
1190#[runtime_builtin(
1191    name = "data.export",
1192    category = "io/data",
1193    summary = "Export dataset to target path.",
1194    keywords = "data,dataset,export",
1195    sink = true,
1196    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1197    descriptor(crate::builtins::io::data::DATA_EXPORT_DESCRIPTOR),
1198    builtin_path = "crate::builtins::io::data"
1199)]
1200async fn data_export_builtin(
1201    path: Value,
1202    format: Value,
1203    target_path: Value,
1204    _rest: Vec<Value>,
1205) -> BuiltinResult<Value> {
1206    let path = parse_string(&path, "data.export path")?;
1207    let format = parse_string(&format, "data.export format")?;
1208    if !format.eq_ignore_ascii_case("data") {
1209        return Err(data_error(
1210            "data.export currently supports only format='data'",
1211        ));
1212    }
1213    let target_path = parse_string(&target_path, "data.export targetPath")?;
1214    copy_dir_recursive(&dataset_root(&path), &dataset_root(&target_path)).await?;
1215    Ok(Value::Bool(true))
1216}
1217
1218#[runtime_builtin(
1219    name = "data.list",
1220    category = "io/data",
1221    summary = "List dataset paths under a prefix.",
1222    keywords = "data,dataset,list",
1223    type_resolver(crate::builtins::io::type_resolvers::data_cell_string_type),
1224    descriptor(crate::builtins::io::data::DATA_LIST_DESCRIPTOR),
1225    builtin_path = "crate::builtins::io::data"
1226)]
1227async fn data_list_builtin(path_prefix: Value, _rest: Vec<Value>) -> BuiltinResult<Value> {
1228    let prefix = parse_string(&path_prefix, "data.list prefix")?;
1229    let root = PathBuf::from(prefix);
1230    let entries = runmat_filesystem::read_dir_async(&root)
1231        .await
1232        .map_err(|err| {
1233            data_error(format!(
1234                "data.list: failed to read '{}': {err}",
1235                root.display()
1236            ))
1237        })?;
1238    let mut values = Vec::new();
1239    for entry in entries {
1240        if !entry.is_dir() {
1241            continue;
1242        }
1243        let candidate = entry.path();
1244        if candidate.extension().and_then(|s| s.to_str()) != Some("data") {
1245            continue;
1246        }
1247        if runmat_filesystem::metadata_async(candidate.join("manifest.json"))
1248            .await
1249            .is_ok()
1250        {
1251            values.push(Value::String(candidate.to_string_lossy().to_string()));
1252        }
1253    }
1254    let cols = values.len();
1255    make_cell(values, 1, cols).map_err(data_error)
1256}
1257
1258#[runtime_builtin(
1259    name = "data.inspect",
1260    category = "io/data",
1261    summary = "Inspect dataset metadata and schema fields.",
1262    keywords = "data,dataset,inspect,schema",
1263    type_resolver(crate::builtins::io::type_resolvers::data_struct_type),
1264    descriptor(crate::builtins::io::data::DATA_INSPECT_DESCRIPTOR),
1265    builtin_path = "crate::builtins::io::data"
1266)]
1267async fn data_inspect_builtin(path: Value) -> BuiltinResult<Value> {
1268    let path = parse_string(&path, "data.inspect path")?;
1269    let root = dataset_root(&path);
1270    let manifest = read_manifest_async(&root).await?;
1271    let mut out = StructValue::new();
1272    out.fields.insert("path".to_string(), Value::String(path));
1273    out.fields
1274        .insert("id".to_string(), Value::String(manifest.dataset_id));
1275    out.fields.insert(
1276        "arrayCount".to_string(),
1277        Value::Num(manifest.arrays.len() as f64),
1278    );
1279    out.fields
1280        .insert("updatedAt".to_string(), Value::String(manifest.updated_at));
1281    Ok(Value::Struct(out))
1282}
1283
1284#[runtime_builtin(
1285    name = "Dataset.path",
1286    category = "io/data",
1287    summary = "Return dataset path.",
1288    keywords = "dataset,path",
1289    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1290    descriptor(crate::builtins::io::data::DATASET_PATH_DESCRIPTOR),
1291    builtin_path = "crate::builtins::io::data"
1292)]
1293async fn dataset_path_builtin(base: Value) -> BuiltinResult<Value> {
1294    let obj = as_object(&base, "Dataset.path")?;
1295    Ok(get_object_prop(obj, "__data_path")?.clone())
1296}
1297
1298#[runtime_builtin(
1299    name = "Dataset.id",
1300    category = "io/data",
1301    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1302    descriptor(crate::builtins::io::data::DATASET_ID_DESCRIPTOR),
1303    builtin_path = "crate::builtins::io::data"
1304)]
1305async fn dataset_id_builtin(base: Value) -> BuiltinResult<Value> {
1306    let obj = as_object(&base, "Dataset.id")?;
1307    Ok(get_object_prop(obj, "__data_id")?.clone())
1308}
1309
1310#[runtime_builtin(
1311    name = "Dataset.version",
1312    category = "io/data",
1313    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1314    descriptor(crate::builtins::io::data::DATASET_VERSION_DESCRIPTOR),
1315    builtin_path = "crate::builtins::io::data"
1316)]
1317async fn dataset_version_builtin(base: Value) -> BuiltinResult<Value> {
1318    let obj = as_object(&base, "Dataset.version")?;
1319    Ok(get_object_prop(obj, "__data_version")?.clone())
1320}
1321
1322#[runtime_builtin(
1323    name = "Dataset.arrays",
1324    category = "io/data",
1325    type_resolver(crate::builtins::io::type_resolvers::data_cell_string_type),
1326    descriptor(crate::builtins::io::data::DATASET_ARRAYS_DESCRIPTOR),
1327    builtin_path = "crate::builtins::io::data"
1328)]
1329async fn dataset_arrays_builtin(base: Value) -> BuiltinResult<Value> {
1330    let path = dataset_path_from_object(&base, "Dataset.arrays")?;
1331    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1332    let values: Vec<Value> = manifest
1333        .arrays
1334        .keys()
1335        .map(|k| Value::String(k.clone()))
1336        .collect();
1337    make_cell(values.clone(), 1, values.len()).map_err(data_error)
1338}
1339
1340#[runtime_builtin(
1341    name = "Dataset.has_array",
1342    category = "io/data",
1343    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1344    descriptor(crate::builtins::io::data::DATASET_HAS_ARRAY_DESCRIPTOR),
1345    builtin_path = "crate::builtins::io::data"
1346)]
1347async fn dataset_has_array_builtin(base: Value, name: Value) -> BuiltinResult<Value> {
1348    let path = dataset_path_from_object(&base, "Dataset.has_array")?;
1349    let name = parse_string(&name, "Dataset.has_array name")?;
1350    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1351    Ok(Value::Bool(manifest.arrays.contains_key(&name)))
1352}
1353
1354#[runtime_builtin(
1355    name = "Dataset.array",
1356    category = "io/data",
1357    type_resolver(crate::builtins::io::type_resolvers::data_array_type),
1358    descriptor(crate::builtins::io::data::DATASET_ARRAY_DESCRIPTOR),
1359    builtin_path = "crate::builtins::io::data"
1360)]
1361async fn dataset_array_builtin(base: Value, name: Value) -> BuiltinResult<Value> {
1362    let path = dataset_path_from_object(&base, "Dataset.array")?;
1363    let name = parse_string(&name, "Dataset.array name")?;
1364    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1365    if !manifest.arrays.contains_key(&name) {
1366        return Err(data_error(format!(
1367            "Dataset.array: array '{name}' not found"
1368        )));
1369    }
1370    Ok(array_object(&path, &name))
1371}
1372
1373#[runtime_builtin(
1374    name = "Dataset.attrs",
1375    category = "io/data",
1376    type_resolver(crate::builtins::io::type_resolvers::data_struct_type),
1377    descriptor(crate::builtins::io::data::DATASET_ATTRS_DESCRIPTOR),
1378    builtin_path = "crate::builtins::io::data"
1379)]
1380async fn dataset_attrs_builtin(base: Value) -> BuiltinResult<Value> {
1381    let path = dataset_path_from_object(&base, "Dataset.attrs")?;
1382    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1383    Ok(attrs_to_struct(&manifest.attrs))
1384}
1385
1386#[runtime_builtin(
1387    name = "Dataset.get_attr",
1388    category = "io/data",
1389    type_resolver(crate::builtins::io::type_resolvers::data_unknown_type),
1390    descriptor(crate::builtins::io::data::DATASET_GET_ATTR_DESCRIPTOR),
1391    builtin_path = "crate::builtins::io::data"
1392)]
1393async fn dataset_get_attr_builtin(
1394    base: Value,
1395    key: Value,
1396    rest: Vec<Value>,
1397) -> BuiltinResult<Value> {
1398    let path = dataset_path_from_object(&base, "Dataset.get_attr")?;
1399    let key = parse_string(&key, "Dataset.get_attr key")?;
1400    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1401    if let Some(value) = manifest.attrs.get(&key) {
1402        return Ok(json_to_value(value));
1403    }
1404    Ok(rest.first().cloned().unwrap_or(Value::Num(0.0)))
1405}
1406
1407#[runtime_builtin(
1408    name = "Dataset.set_attr",
1409    category = "io/data",
1410    sink = true,
1411    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1412    descriptor(crate::builtins::io::data::DATASET_SET_ATTR_DESCRIPTOR),
1413    builtin_path = "crate::builtins::io::data"
1414)]
1415async fn dataset_set_attr_builtin(base: Value, key: Value, value: Value) -> BuiltinResult<Value> {
1416    let path = dataset_path_from_object(&base, "Dataset.set_attr")?;
1417    let key = parse_string(&key, "Dataset.set_attr key")?;
1418    let root = dataset_root(&path);
1419    let mut manifest = read_manifest_async(&root).await?;
1420    manifest.attrs.insert(key, value_to_json(&value));
1421    manifest.updated_at = now_rfc3339();
1422    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1423    write_manifest_async(&root, &manifest).await?;
1424    Ok(Value::Bool(true))
1425}
1426
1427#[runtime_builtin(
1428    name = "Dataset.set_attrs",
1429    category = "io/data",
1430    sink = true,
1431    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1432    descriptor(crate::builtins::io::data::DATASET_SET_ATTRS_DESCRIPTOR),
1433    builtin_path = "crate::builtins::io::data"
1434)]
1435async fn dataset_set_attrs_builtin(base: Value, attrs: Value) -> BuiltinResult<Value> {
1436    let path = dataset_path_from_object(&base, "Dataset.set_attrs")?;
1437    let Value::Struct(incoming) = attrs else {
1438        return Err(data_error("Dataset.set_attrs: attrs must be a struct"));
1439    };
1440    let root = dataset_root(&path);
1441    let mut manifest = read_manifest_async(&root).await?;
1442    for (k, v) in incoming.fields {
1443        manifest.attrs.insert(k, value_to_json(&v));
1444    }
1445    manifest.updated_at = now_rfc3339();
1446    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1447    write_manifest_async(&root, &manifest).await?;
1448    Ok(Value::Bool(true))
1449}
1450
1451#[runtime_builtin(
1452    name = "Dataset.begin",
1453    category = "io/data",
1454    type_resolver(crate::builtins::io::type_resolvers::data_tx_type),
1455    descriptor(crate::builtins::io::data::DATASET_BEGIN_DESCRIPTOR),
1456    builtin_path = "crate::builtins::io::data"
1457)]
1458async fn dataset_begin_builtin(base: Value, _rest: Vec<Value>) -> BuiltinResult<Value> {
1459    let path = dataset_path_from_object(&base, "Dataset.begin")?;
1460    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1461    let tx_id = start_tx(path.clone(), manifest.txn_sequence)?;
1462    tracing::info!(
1463        target: "runmat.data",
1464        dataset = path,
1465        tx_id = tx_id,
1466        base_sequence = manifest.txn_sequence,
1467        "data transaction begin"
1468    );
1469    Ok(transaction_object(&path, &tx_id))
1470}
1471
1472#[runtime_builtin(
1473    name = "Dataset.snapshot",
1474    category = "io/data",
1475    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1476    descriptor(crate::builtins::io::data::DATASET_SNAPSHOT_DESCRIPTOR),
1477    builtin_path = "crate::builtins::io::data"
1478)]
1479async fn dataset_snapshot_builtin(
1480    base: Value,
1481    label: Value,
1482    _rest: Vec<Value>,
1483) -> BuiltinResult<Value> {
1484    let path = dataset_path_from_object(&base, "Dataset.snapshot")?;
1485    let label = parse_string(&label, "Dataset.snapshot label")?;
1486    let root = dataset_root(&path);
1487    let snapshots = root.join(".snapshots");
1488    runmat_filesystem::create_dir_all_async(&snapshots)
1489        .await
1490        .map_err(|err| {
1491            data_error(format!(
1492                "Dataset.snapshot: failed to create snapshots dir: {err}"
1493            ))
1494        })?;
1495    let src = manifest_path(&root);
1496    let dst = snapshots.join(format!("{}.manifest.json", sanitize_label(&label)));
1497    copy_file(&src, &dst).await?;
1498    Ok(Value::String(dst.to_string_lossy().to_string()))
1499}
1500
1501#[runtime_builtin(
1502    name = "Dataset.refresh",
1503    category = "io/data",
1504    type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1505    descriptor(crate::builtins::io::data::DATASET_REFRESH_DESCRIPTOR),
1506    builtin_path = "crate::builtins::io::data"
1507)]
1508async fn dataset_refresh_builtin(base: Value) -> BuiltinResult<Value> {
1509    let path = dataset_path_from_object(&base, "Dataset.refresh")?;
1510    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1511    let mut ds = dataset_object(&path, &manifest);
1512    hydrate_dataset_descriptor_async(&path, &mut ds).await;
1513    Ok(ds)
1514}
1515
1516#[runtime_builtin(
1517    name = "DataArray.name",
1518    category = "io/data",
1519    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1520    descriptor(crate::builtins::io::data::DATAARRAY_NAME_DESCRIPTOR),
1521    builtin_path = "crate::builtins::io::data"
1522)]
1523async fn data_array_name_builtin(base: Value) -> BuiltinResult<Value> {
1524    let obj = as_object(&base, "DataArray.name")?;
1525    Ok(get_object_prop(obj, "__array_name")?.clone())
1526}
1527
1528#[runtime_builtin(
1529    name = "DataArray.dtype",
1530    category = "io/data",
1531    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1532    descriptor(crate::builtins::io::data::DATAARRAY_DTYPE_DESCRIPTOR),
1533    builtin_path = "crate::builtins::io::data"
1534)]
1535async fn data_array_dtype_builtin(base: Value) -> BuiltinResult<Value> {
1536    let (path, name) = array_identity(&base, "DataArray.dtype")?;
1537    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1538    let meta = manifest
1539        .arrays
1540        .get(&name)
1541        .ok_or_else(|| data_error(format!("DataArray.dtype: array '{name}' not found")))?;
1542    Ok(Value::String(meta.dtype.clone()))
1543}
1544
1545#[runtime_builtin(
1546    name = "DataArray.shape",
1547    category = "io/data",
1548    type_resolver(crate::builtins::io::type_resolvers::data_shape_tensor_type),
1549    descriptor(crate::builtins::io::data::DATAARRAY_SHAPE_DESCRIPTOR),
1550    builtin_path = "crate::builtins::io::data"
1551)]
1552async fn data_array_shape_builtin(base: Value) -> BuiltinResult<Value> {
1553    let (path, name) = array_identity(&base, "DataArray.shape")?;
1554    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1555    let meta = manifest
1556        .arrays
1557        .get(&name)
1558        .ok_or_else(|| data_error(format!("DataArray.shape: array '{name}' not found")))?;
1559    let values = meta.shape.iter().map(|v| *v as f64).collect::<Vec<_>>();
1560    let tensor = Tensor::new(values, vec![1, meta.shape.len()])
1561        .map_err(|err| data_error(format!("DataArray.shape: {err}")))?;
1562    Ok(Value::Tensor(tensor))
1563}
1564
1565#[runtime_builtin(
1566    name = "DataArray.rank",
1567    category = "io/data",
1568    type_resolver(crate::builtins::io::type_resolvers::data_int_type),
1569    descriptor(crate::builtins::io::data::DATAARRAY_RANK_DESCRIPTOR),
1570    builtin_path = "crate::builtins::io::data"
1571)]
1572async fn data_array_rank_builtin(base: Value) -> BuiltinResult<Value> {
1573    let (path, name) = array_identity(&base, "DataArray.rank")?;
1574    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1575    let meta = manifest
1576        .arrays
1577        .get(&name)
1578        .ok_or_else(|| data_error(format!("DataArray.rank: array '{name}' not found")))?;
1579    Ok(Value::Num(meta.shape.len() as f64))
1580}
1581
1582#[runtime_builtin(
1583    name = "DataArray.chunk_shape",
1584    category = "io/data",
1585    type_resolver(crate::builtins::io::type_resolvers::data_shape_tensor_type),
1586    descriptor(crate::builtins::io::data::DATAARRAY_CHUNK_SHAPE_DESCRIPTOR),
1587    builtin_path = "crate::builtins::io::data"
1588)]
1589async fn data_array_chunk_shape_builtin(base: Value) -> BuiltinResult<Value> {
1590    let (path, name) = array_identity(&base, "DataArray.chunk_shape")?;
1591    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1592    let meta = manifest
1593        .arrays
1594        .get(&name)
1595        .ok_or_else(|| data_error(format!("DataArray.chunk_shape: array '{name}' not found")))?;
1596    let values = meta
1597        .chunk_shape
1598        .iter()
1599        .map(|v| *v as f64)
1600        .collect::<Vec<_>>();
1601    let tensor = Tensor::new(values, vec![1, meta.chunk_shape.len()])
1602        .map_err(|err| data_error(format!("DataArray.chunk_shape: {err}")))?;
1603    Ok(Value::Tensor(tensor))
1604}
1605
1606#[runtime_builtin(
1607    name = "DataArray.codec",
1608    category = "io/data",
1609    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1610    descriptor(crate::builtins::io::data::DATAARRAY_CODEC_DESCRIPTOR),
1611    builtin_path = "crate::builtins::io::data"
1612)]
1613async fn data_array_codec_builtin(base: Value) -> BuiltinResult<Value> {
1614    let (path, name) = array_identity(&base, "DataArray.codec")?;
1615    let manifest = read_manifest_async(&dataset_root(&path)).await?;
1616    let meta = manifest
1617        .arrays
1618        .get(&name)
1619        .ok_or_else(|| data_error(format!("DataArray.codec: array '{name}' not found")))?;
1620    Ok(Value::String(meta.codec.clone()))
1621}
1622
1623#[runtime_builtin(
1624    name = "DataArray.read",
1625    category = "io/data",
1626    type_resolver(crate::builtins::io::type_resolvers::data_tensor_type),
1627    descriptor(crate::builtins::io::data::DATAARRAY_READ_DESCRIPTOR),
1628    builtin_path = "crate::builtins::io::data"
1629)]
1630async fn data_array_read_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
1631    let (path, name) = array_identity(&base, "DataArray.read")?;
1632    let root = dataset_root(&path);
1633    let manifest = read_manifest_async(&root).await?;
1634    let meta = manifest
1635        .arrays
1636        .get(&name)
1637        .ok_or_else(|| data_error(format!("DataArray.read: array '{name}' not found")))?;
1638    let payload = read_array_payload_async(&root, meta).await?;
1639    let sliced = if let Some(slice_spec) = rest.first() {
1640        read_slice_payload(&payload, slice_spec)?
1641    } else {
1642        payload
1643    };
1644    sliced
1645        .into_value()
1646        .map_err(|err| data_error(format!("DataArray.read: {err}")))
1647}
1648
1649#[runtime_builtin(
1650    name = "DataArray.write",
1651    category = "io/data",
1652    sink = true,
1653    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1654    descriptor(crate::builtins::io::data::DATAARRAY_WRITE_DESCRIPTOR),
1655    builtin_path = "crate::builtins::io::data"
1656)]
1657async fn data_array_write_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
1658    let (path, name) = array_identity(&base, "DataArray.write")?;
1659    let (slice_spec, value) = match rest.as_slice() {
1660        [v] => (None, v),
1661        [slice, v] => (Some(slice), v),
1662        _ => {
1663            return Err(data_error(
1664                "DataArray.write expects values or (sliceSpec, values) arguments",
1665            ))
1666        }
1667    };
1668    write_array_full_async(&path, &name, slice_spec, value).await?;
1669    Ok(Value::Bool(true))
1670}
1671
1672#[runtime_builtin(
1673    name = "DataArray.resize",
1674    category = "io/data",
1675    sink = true,
1676    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1677    descriptor(crate::builtins::io::data::DATAARRAY_RESIZE_DESCRIPTOR),
1678    builtin_path = "crate::builtins::io::data"
1679)]
1680async fn data_array_resize_builtin(
1681    base: Value,
1682    new_shape: Value,
1683    _rest: Vec<Value>,
1684) -> BuiltinResult<Value> {
1685    let (path, name) = array_identity(&base, "DataArray.resize")?;
1686    let shape = parse_shape_from_value(&new_shape)?;
1687    let root = dataset_root(&path);
1688    let mut manifest = read_manifest_async(&root).await?;
1689    let meta = manifest
1690        .arrays
1691        .get_mut(&name)
1692        .ok_or_else(|| data_error(format!("DataArray.resize: array '{name}' not found")))?;
1693    meta.shape = shape.clone();
1694    let payload = DataArrayPayload::zeros(meta.dtype.clone(), shape.clone());
1695    let (payload_path, chunk_index_path) =
1696        write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
1697    meta.data_path = make_rel_data_path(&root, &payload_path)?;
1698    meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
1699    manifest.updated_at = now_rfc3339();
1700    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1701    write_manifest_async(&root, &manifest).await?;
1702    Ok(Value::Bool(true))
1703}
1704
1705#[runtime_builtin(
1706    name = "DataArray.fill",
1707    category = "io/data",
1708    sink = true,
1709    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1710    descriptor(crate::builtins::io::data::DATAARRAY_FILL_DESCRIPTOR),
1711    builtin_path = "crate::builtins::io::data"
1712)]
1713async fn data_array_fill_builtin(
1714    base: Value,
1715    value: Value,
1716    _rest: Vec<Value>,
1717) -> BuiltinResult<Value> {
1718    let (path, name) = array_identity(&base, "DataArray.fill")?;
1719    let root = dataset_root(&path);
1720    let mut manifest = read_manifest_async(&root).await?;
1721    let meta = manifest
1722        .arrays
1723        .get_mut(&name)
1724        .ok_or_else(|| data_error(format!("DataArray.fill: array '{name}' not found")))?;
1725    let payload = DataArrayPayload::filled(meta.dtype.clone(), meta.shape.clone(), &value)?;
1726    let (payload_path, chunk_index_path) =
1727        write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
1728    meta.data_path = make_rel_data_path(&root, &payload_path)?;
1729    meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
1730    manifest.updated_at = now_rfc3339();
1731    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1732    write_manifest_async(&root, &manifest).await?;
1733    Ok(Value::Bool(true))
1734}
1735
1736#[runtime_builtin(
1737    name = "DataTransaction.id",
1738    category = "io/data",
1739    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1740    descriptor(crate::builtins::io::data::DATATX_ID_DESCRIPTOR),
1741    builtin_path = "crate::builtins::io::data"
1742)]
1743async fn data_tx_id_builtin(base: Value) -> BuiltinResult<Value> {
1744    let obj = as_object(&base, "DataTransaction.id")?;
1745    Ok(get_object_prop(obj, "__tx_id")?.clone())
1746}
1747
1748#[runtime_builtin(
1749    name = "DataTransaction.write",
1750    category = "io/data",
1751    sink = true,
1752    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1753    descriptor(crate::builtins::io::data::DATATX_WRITE_DESCRIPTOR),
1754    builtin_path = "crate::builtins::io::data"
1755)]
1756async fn data_tx_write_builtin(
1757    base: Value,
1758    array_name: Value,
1759    slice: Value,
1760    values: Value,
1761    _rest: Vec<Value>,
1762) -> BuiltinResult<Value> {
1763    let tx_id = tx_id_from_object(&base, "DataTransaction.write")?;
1764    let array_name = parse_string(&array_name, "DataTransaction.write arrayName")?;
1765    with_tx_mut(&tx_id, |tx| {
1766        if tx.status != TxnStatus::Open {
1767            return Err(data_error("DataTransaction.write: transaction is not open"));
1768        }
1769        tx.writes.push(PendingWrite {
1770            array: array_name,
1771            slice_spec: Some(slice),
1772            value: values,
1773        });
1774        Ok(())
1775    })?;
1776    Ok(Value::Bool(true))
1777}
1778
1779#[runtime_builtin(
1780    name = "DataTransaction.set_attr",
1781    category = "io/data",
1782    sink = true,
1783    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1784    descriptor(crate::builtins::io::data::DATATX_SET_ATTR_DESCRIPTOR),
1785    builtin_path = "crate::builtins::io::data"
1786)]
1787async fn data_tx_set_attr_builtin(base: Value, key: Value, value: Value) -> BuiltinResult<Value> {
1788    let tx_id = tx_id_from_object(&base, "DataTransaction.set_attr")?;
1789    let key = parse_string(&key, "DataTransaction.set_attr key")?;
1790    with_tx_mut(&tx_id, |tx| {
1791        if tx.status != TxnStatus::Open {
1792            return Err(data_error(
1793                "DataTransaction.set_attr: transaction is not open",
1794            ));
1795        }
1796        tx.attrs.insert(key, value);
1797        Ok(())
1798    })?;
1799    Ok(Value::Bool(true))
1800}
1801
1802#[runtime_builtin(
1803    name = "DataTransaction.set_attrs",
1804    category = "io/data",
1805    sink = true,
1806    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1807    descriptor(crate::builtins::io::data::DATATX_SET_ATTRS_DESCRIPTOR),
1808    builtin_path = "crate::builtins::io::data"
1809)]
1810async fn data_tx_set_attrs_builtin(base: Value, attrs: Value) -> BuiltinResult<Value> {
1811    let tx_id = tx_id_from_object(&base, "DataTransaction.set_attrs")?;
1812    let Value::Struct(incoming) = attrs else {
1813        return Err(data_error(
1814            "DataTransaction.set_attrs: attrs must be struct",
1815        ));
1816    };
1817    with_tx_mut(&tx_id, |tx| {
1818        if tx.status != TxnStatus::Open {
1819            return Err(data_error(
1820                "DataTransaction.set_attrs: transaction is not open",
1821            ));
1822        }
1823        for (k, v) in incoming.fields {
1824            tx.attrs.insert(k, v);
1825        }
1826        Ok(())
1827    })?;
1828    Ok(Value::Bool(true))
1829}
1830
1831#[runtime_builtin(
1832    name = "DataTransaction.resize",
1833    category = "io/data",
1834    sink = true,
1835    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1836    descriptor(crate::builtins::io::data::DATATX_RESIZE_DESCRIPTOR),
1837    builtin_path = "crate::builtins::io::data"
1838)]
1839async fn data_tx_resize_builtin(
1840    base: Value,
1841    array_name: Value,
1842    new_shape: Value,
1843    _rest: Vec<Value>,
1844) -> BuiltinResult<Value> {
1845    let tx_id = tx_id_from_object(&base, "DataTransaction.resize")?;
1846    let array_name = parse_string(&array_name, "DataTransaction.resize arrayName")?;
1847    let shape = parse_shape_from_value(&new_shape)?;
1848    with_tx_mut(&tx_id, |tx| {
1849        if tx.status != TxnStatus::Open {
1850            return Err(data_error(
1851                "DataTransaction.resize: transaction is not open",
1852            ));
1853        }
1854        tx.resizes.push(PendingResize {
1855            array: array_name,
1856            shape,
1857        });
1858        Ok(())
1859    })?;
1860    Ok(Value::Bool(true))
1861}
1862
1863#[runtime_builtin(
1864    name = "DataTransaction.fill",
1865    category = "io/data",
1866    sink = true,
1867    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1868    descriptor(crate::builtins::io::data::DATATX_FILL_DESCRIPTOR),
1869    builtin_path = "crate::builtins::io::data"
1870)]
1871async fn data_tx_fill_builtin(
1872    base: Value,
1873    array_name: Value,
1874    value: Value,
1875    rest: Vec<Value>,
1876) -> BuiltinResult<Value> {
1877    let tx_id = tx_id_from_object(&base, "DataTransaction.fill")?;
1878    let array_name = parse_string(&array_name, "DataTransaction.fill arrayName")?;
1879    let slice_spec = rest.first().cloned();
1880    with_tx_mut(&tx_id, |tx| {
1881        if tx.status != TxnStatus::Open {
1882            return Err(data_error("DataTransaction.fill: transaction is not open"));
1883        }
1884        tx.fills.push(PendingFill {
1885            array: array_name,
1886            slice_spec,
1887            value,
1888        });
1889        Ok(())
1890    })?;
1891    Ok(Value::Bool(true))
1892}
1893
1894#[runtime_builtin(
1895    name = "DataTransaction.delete_array",
1896    category = "io/data",
1897    sink = true,
1898    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1899    descriptor(crate::builtins::io::data::DATATX_DELETE_ARRAY_DESCRIPTOR),
1900    builtin_path = "crate::builtins::io::data"
1901)]
1902async fn data_tx_delete_array_builtin(base: Value, array_name: Value) -> BuiltinResult<Value> {
1903    let tx_id = tx_id_from_object(&base, "DataTransaction.delete_array")?;
1904    let array_name = parse_string(&array_name, "DataTransaction.delete_array arrayName")?;
1905    with_tx_mut(&tx_id, |tx| {
1906        if tx.status != TxnStatus::Open {
1907            return Err(data_error(
1908                "DataTransaction.delete_array: transaction is not open",
1909            ));
1910        }
1911        tx.delete_arrays.push(array_name);
1912        Ok(())
1913    })?;
1914    Ok(Value::Bool(true))
1915}
1916
1917#[runtime_builtin(
1918    name = "DataTransaction.create_array",
1919    category = "io/data",
1920    sink = true,
1921    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1922    descriptor(crate::builtins::io::data::DATATX_CREATE_ARRAY_DESCRIPTOR),
1923    builtin_path = "crate::builtins::io::data"
1924)]
1925async fn data_tx_create_array_builtin(
1926    base: Value,
1927    array_name: Value,
1928    meta: Value,
1929) -> BuiltinResult<Value> {
1930    let tx_id = tx_id_from_object(&base, "DataTransaction.create_array")?;
1931    let array_name = parse_string(&array_name, "DataTransaction.create_array arrayName")?;
1932    let meta = parse_array_meta(&array_name, &meta)?;
1933    with_tx_mut(&tx_id, |tx| {
1934        if tx.status != TxnStatus::Open {
1935            return Err(data_error(
1936                "DataTransaction.create_array: transaction is not open",
1937            ));
1938        }
1939        tx.create_arrays.push(PendingCreateArray {
1940            array: array_name,
1941            meta,
1942        });
1943        Ok(())
1944    })?;
1945    Ok(Value::Bool(true))
1946}
1947
1948#[runtime_builtin(
1949    name = "DataTransaction.commit",
1950    category = "io/data",
1951    sink = true,
1952    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1953    descriptor(crate::builtins::io::data::DATATX_COMMIT_DESCRIPTOR),
1954    builtin_path = "crate::builtins::io::data"
1955)]
1956async fn data_tx_commit_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
1957    let tx_id = tx_id_from_object(&base, "DataTransaction.commit")?;
1958    let (dataset_path, base_sequence, writes, resizes, fills, create_arrays, delete_arrays, attrs) =
1959        with_tx(&tx_id, |tx| {
1960            if tx.status != TxnStatus::Open {
1961                return Err(data_error(
1962                    "DataTransaction.commit: transaction is not open",
1963                ));
1964            }
1965            Ok((
1966                tx.dataset_path.clone(),
1967                tx.base_sequence,
1968                tx.writes.clone(),
1969                tx.resizes.clone(),
1970                tx.fills.clone(),
1971                tx.create_arrays.clone(),
1972                tx.delete_arrays.clone(),
1973                tx.attrs.clone(),
1974            ))
1975        })?;
1976
1977    let write_ops = writes.len();
1978    let resize_ops = resizes.len();
1979    let fill_ops = fills.len();
1980    let create_ops = create_arrays.len();
1981    let delete_ops = delete_arrays.len();
1982    let attr_updates = attrs.len();
1983
1984    let root = dataset_root(&dataset_path);
1985    let mut manifest = read_manifest_async(&root).await?;
1986    ensure_manifest_sequence(base_sequence, &manifest)?;
1987    if let Some(Value::Struct(options)) = rest.first() {
1988        if let Some(expected) = options.fields.get("if_manifest") {
1989            let expected = parse_string(expected, "DataTransaction.commit if_manifest")?;
1990            let actual = manifest_version_token(&manifest);
1991            if expected != actual {
1992                tracing::warn!(
1993                    target: "runmat.data",
1994                    tx_id = tx_id,
1995                    expected_manifest = expected,
1996                    actual_manifest = actual,
1997                    "data transaction manifest conflict"
1998                );
1999                return Err(data_error(
2000                    "MANIFEST_CONFLICT: if_manifest precondition failed",
2001                ));
2002            }
2003        }
2004    }
2005    for create in create_arrays {
2006        create_array_in_manifest(&root, &mut manifest, &create.array, create.meta).await?;
2007    }
2008    for resize in resizes {
2009        resize_array_in_manifest(&root, &mut manifest, &resize.array, resize.shape).await?;
2010    }
2011    for fill in fills {
2012        fill_array_in_manifest(
2013            &root,
2014            &mut manifest,
2015            &fill.array,
2016            fill.slice_spec.as_ref(),
2017            &fill.value,
2018        )
2019        .await?;
2020    }
2021    for write in writes {
2022        apply_write_to_manifest_async(
2023            &root,
2024            &mut manifest,
2025            &write.array,
2026            write.slice_spec.as_ref(),
2027            &write.value,
2028        )
2029        .await?;
2030    }
2031    for array_name in delete_arrays {
2032        delete_array_in_manifest_async(&root, &mut manifest, &array_name).await?;
2033    }
2034    for (k, v) in attrs {
2035        manifest.attrs.insert(k, value_to_json(&v));
2036    }
2037    manifest.updated_at = now_rfc3339();
2038    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2039    write_manifest_async(&root, &manifest).await?;
2040    with_tx_mut(&tx_id, |tx| {
2041        tx.status = TxnStatus::Committed;
2042        Ok(())
2043    })?;
2044    tracing::info!(
2045        target: "runmat.data",
2046        dataset = dataset_path,
2047        tx_id = tx_id,
2048        write_ops = write_ops,
2049        resize_ops = resize_ops,
2050        fill_ops = fill_ops,
2051        create_ops = create_ops,
2052        delete_ops = delete_ops,
2053        attr_updates = attr_updates,
2054        next_sequence = manifest.txn_sequence,
2055        "data transaction commit"
2056    );
2057    remove_tx(&tx_id)?;
2058    Ok(Value::Bool(true))
2059}
2060
2061#[runtime_builtin(
2062    name = "commit",
2063    category = "io/data",
2064    sink = true,
2065    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2066    descriptor(crate::builtins::io::data::COMMIT_ALIAS_DESCRIPTOR),
2067    builtin_path = "crate::builtins::io::data"
2068)]
2069async fn data_tx_commit_alias_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
2070    match &base {
2071        Value::Object(obj) if obj.class_name == "DataTransaction" => {
2072            data_tx_commit_builtin(base, rest).await
2073        }
2074        Value::HandleObject(handle) if handle.class_name == "DataTransaction" => {
2075            data_tx_commit_builtin(base, rest).await
2076        }
2077        _ => Err(data_error(
2078            "commit: receiver must be a DataTransaction (use tx = ds.begin())",
2079        )),
2080    }
2081}
2082
2083#[runtime_builtin(
2084    name = "DataTransaction.abort",
2085    category = "io/data",
2086    sink = true,
2087    type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2088    descriptor(crate::builtins::io::data::DATATX_ABORT_DESCRIPTOR),
2089    builtin_path = "crate::builtins::io::data"
2090)]
2091async fn data_tx_abort_builtin(base: Value) -> BuiltinResult<Value> {
2092    let tx_id = tx_id_from_object(&base, "DataTransaction.abort")?;
2093    with_tx_mut(&tx_id, |tx| {
2094        tx.status = TxnStatus::Aborted;
2095        Ok(())
2096    })?;
2097    tracing::info!(
2098        target: "runmat.data",
2099        tx_id = tx_id,
2100        "data transaction abort"
2101    );
2102    remove_tx(&tx_id)?;
2103    Ok(Value::Bool(true))
2104}
2105
2106#[runtime_builtin(
2107    name = "DataTransaction.status",
2108    category = "io/data",
2109    type_resolver(crate::builtins::io::type_resolvers::data_string_type),
2110    descriptor(crate::builtins::io::data::DATATX_STATUS_DESCRIPTOR),
2111    builtin_path = "crate::builtins::io::data"
2112)]
2113async fn data_tx_status_builtin(base: Value) -> BuiltinResult<Value> {
2114    let tx_id = tx_id_from_object(&base, "DataTransaction.status")?;
2115    with_tx(&tx_id, |tx| {
2116        let status = match tx.status {
2117            TxnStatus::Open => "open",
2118            TxnStatus::Committed => "committed",
2119            TxnStatus::Aborted => "aborted",
2120        };
2121        Ok(Value::String(status.to_string()))
2122    })
2123}
2124
2125fn dataset_path_from_object(base: &Value, context: &str) -> BuiltinResult<String> {
2126    let obj = as_object(base, context)?;
2127    parse_string(get_object_prop(obj, "__data_path")?, context)
2128}
2129
2130fn tx_id_from_object(base: &Value, context: &str) -> BuiltinResult<String> {
2131    let obj = as_object(base, context)?;
2132    parse_string(get_object_prop(obj, "__tx_id")?, context)
2133}
2134
2135fn array_identity(base: &Value, context: &str) -> BuiltinResult<(String, String)> {
2136    let obj = as_object(base, context)?;
2137    let path = parse_string(get_object_prop(obj, "__data_path")?, context)?;
2138    let name = parse_string(get_object_prop(obj, "__array_name")?, context)?;
2139    Ok((path, name))
2140}
2141
2142fn as_object<'a>(value: &'a Value, context: &str) -> BuiltinResult<&'a ObjectInstance> {
2143    match value {
2144        Value::Object(obj) => Ok(obj),
2145        _ => Err(data_error(format!("{context}: expected object receiver"))),
2146    }
2147}
2148
2149async fn hydrate_dataset_descriptor_async(path: &str, dataset: &mut Value) {
2150    let request = runmat_filesystem::data_contract::DataManifestRequest {
2151        path: path.to_string(),
2152        version: None,
2153    };
2154    let descriptor = match runmat_filesystem::data_manifest_descriptor_async(&request).await {
2155        Ok(descriptor) => descriptor,
2156        Err(_) => return,
2157    };
2158    let Value::Object(obj) = dataset else {
2159        return;
2160    };
2161    if !descriptor.dataset_id.is_empty() {
2162        obj.properties.insert(
2163            "__data_id".to_string(),
2164            Value::String(descriptor.dataset_id),
2165        );
2166    }
2167    obj.properties.insert(
2168        "__data_version".to_string(),
2169        Value::String(format!(
2170            "{}:{}",
2171            descriptor.updated_at, descriptor.txn_sequence
2172        )),
2173    );
2174}
2175
2176fn sanitize_label(label: &str) -> String {
2177    label
2178        .chars()
2179        .map(|ch| {
2180            if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' {
2181                ch
2182            } else {
2183                '_'
2184            }
2185        })
2186        .collect()
2187}
2188
2189async fn copy_file(src: &PathBuf, dst: &PathBuf) -> BuiltinResult<()> {
2190    let bytes = runmat_filesystem::read_async(src)
2191        .await
2192        .map_err(|err| data_error(format!("failed to open '{}': {err}", src.display())))?;
2193    let parent = dst.parent().ok_or_else(|| {
2194        data_error(format!(
2195            "invalid destination path '{}': missing parent",
2196            dst.display()
2197        ))
2198    })?;
2199    runmat_filesystem::create_dir_all_async(parent)
2200        .await
2201        .map_err(|err| data_error(format!("failed to create '{}': {err}", parent.display())))?;
2202    runmat_filesystem::write_async(dst, &bytes)
2203        .await
2204        .map_err(|err| {
2205            data_error(format!(
2206                "failed to copy '{}' -> '{}': {err}",
2207                src.display(),
2208                dst.display()
2209            ))
2210        })?;
2211    Ok(())
2212}
2213
2214fn make_rel_data_path(
2215    root: &std::path::Path,
2216    payload_path: &std::path::Path,
2217) -> BuiltinResult<String> {
2218    let rel = payload_path
2219        .strip_prefix(root)
2220        .map_err(|err| data_error(format!("failed to compute relative data path: {err}")))?;
2221    Ok(rel.to_string_lossy().to_string())
2222}
2223
2224async fn create_array_in_manifest(
2225    root: &std::path::Path,
2226    manifest: &mut DataManifest,
2227    array_name: &str,
2228    mut meta: DataArrayMeta,
2229) -> BuiltinResult<()> {
2230    if manifest.arrays.contains_key(array_name) {
2231        return Err(data_error(format!(
2232            "DataTransaction.create_array: array '{array_name}' already exists"
2233        )));
2234    }
2235    let payload = DataArrayPayload::zeros(meta.dtype.clone(), meta.shape.clone());
2236    let (payload_path, chunk_index_path) =
2237        write_array_payload_async(root, array_name, &payload, &meta.chunk_shape).await?;
2238    meta.data_path = make_rel_data_path(root, &payload_path)?;
2239    meta.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2240    manifest.arrays.insert(array_name.to_string(), meta);
2241    Ok(())
2242}
2243
2244async fn resize_array_in_manifest(
2245    root: &std::path::Path,
2246    manifest: &mut DataManifest,
2247    array_name: &str,
2248    shape: Vec<usize>,
2249) -> BuiltinResult<()> {
2250    let meta = manifest
2251        .arrays
2252        .get_mut(array_name)
2253        .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2254    meta.shape = shape.clone();
2255    let payload = DataArrayPayload::zeros(meta.dtype.clone(), shape.clone());
2256    let (payload_path, chunk_index_path) =
2257        write_array_payload_async(root, array_name, &payload, &meta.chunk_shape).await?;
2258    meta.data_path = make_rel_data_path(root, &payload_path)?;
2259    meta.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2260    Ok(())
2261}
2262
2263async fn fill_array_in_manifest(
2264    root: &std::path::Path,
2265    manifest: &mut DataManifest,
2266    array_name: &str,
2267    slice_spec: Option<&Value>,
2268    value: &Value,
2269) -> BuiltinResult<()> {
2270    let meta: DataArrayMeta = manifest
2271        .arrays
2272        .get(array_name)
2273        .cloned()
2274        .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2275    let payload = read_array_payload_async(root, &meta).await?;
2276    let next_payload = if let Some(slice_spec) = slice_spec {
2277        let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
2278        let target_shape: Vec<usize> = ranges
2279            .iter()
2280            .map(|r| r.end.saturating_sub(r.start))
2281            .collect();
2282        let rhs = DataArrayPayload::filled(payload.dtype.clone(), target_shape, value)?
2283            .into_value()
2284            .map_err(|err| data_error(format!("DataTransaction.fill: {err}")))?;
2285        write_slice_payload(&payload, slice_spec, &rhs)?
2286    } else {
2287        DataArrayPayload::filled(payload.dtype.clone(), payload.shape.clone(), value)?
2288    };
2289    let (payload_path, chunk_index_path) =
2290        write_array_payload_async(root, array_name, &next_payload, &meta.chunk_shape).await?;
2291    if let Some(updated) = manifest.arrays.get_mut(array_name) {
2292        updated.shape = next_payload.shape.clone();
2293        updated.data_path = make_rel_data_path(root, &payload_path)?;
2294        updated.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2295    }
2296    Ok(())
2297}
2298
2299async fn delete_array_in_manifest_async(
2300    root: &std::path::Path,
2301    manifest: &mut DataManifest,
2302    array_name: &str,
2303) -> BuiltinResult<()> {
2304    let removed = manifest.arrays.remove(array_name);
2305    if removed.is_none() {
2306        return Err(data_error(format!(
2307            "DataTransaction.delete_array: array '{array_name}' not found"
2308        )));
2309    }
2310    let array_dir = root.join("arrays").join(array_name);
2311    if runmat_filesystem::metadata_async(&array_dir).await.is_ok() {
2312        runmat_filesystem::remove_dir_all_async(&array_dir)
2313            .await
2314            .map_err(|err| {
2315                data_error(format!(
2316                    "DataTransaction.delete_array: failed to remove '{}': {err}",
2317                    array_dir.display()
2318                ))
2319            })?;
2320    }
2321    Ok(())
2322}
2323
2324fn parse_array_meta(array_name: &str, meta: &Value) -> BuiltinResult<DataArrayMeta> {
2325    let Value::Struct(meta_struct) = meta else {
2326        return Err(data_error(
2327            "DataTransaction.create_array: meta must be a struct",
2328        ));
2329    };
2330    let dtype = meta_struct
2331        .fields
2332        .get("dtype")
2333        .map(|v| parse_string(v, "DataTransaction.create_array dtype"))
2334        .transpose()?
2335        .unwrap_or_else(|| "f64".to_string());
2336    let shape = meta_struct
2337        .fields
2338        .get("shape")
2339        .map(parse_shape_from_value)
2340        .transpose()?
2341        .unwrap_or_else(|| vec![0, 0]);
2342    let chunk_shape = meta_struct
2343        .fields
2344        .get("chunk")
2345        .map(parse_shape_from_value)
2346        .transpose()?
2347        .unwrap_or_else(|| default_chunk_shape(&shape));
2348    let codec = meta_struct
2349        .fields
2350        .get("codec")
2351        .map(|v| parse_string(v, "DataTransaction.create_array codec"))
2352        .transpose()?
2353        .unwrap_or_else(|| "zstd".to_string());
2354    Ok(DataArrayMeta {
2355        dtype,
2356        shape,
2357        chunk_shape,
2358        order: "column_major".to_string(),
2359        codec,
2360        chunk_index_path: Some(format!("arrays/{array_name}/chunks/index.json")),
2361        data_path: format!("arrays/{array_name}/data.f64.json"),
2362    })
2363}
2364
2365fn default_chunk_shape(shape: &[usize]) -> Vec<usize> {
2366    if shape.is_empty() {
2367        return vec![1024];
2368    }
2369    let mut out = shape.to_vec();
2370    if out.len() == 1 {
2371        out[0] = out[0].clamp(1, 65_536);
2372        return out;
2373    }
2374    out[0] = out[0].clamp(1, 256);
2375    out[1] = out[1].clamp(1, 256);
2376    for dim in out.iter_mut().skip(2) {
2377        *dim = (*dim).clamp(1, 8);
2378    }
2379    out
2380}
2381
2382#[async_recursion::async_recursion(?Send)]
2383async fn copy_dir_recursive(src: &PathBuf, dst: &PathBuf) -> BuiltinResult<()> {
2384    let metadata = runmat_filesystem::metadata_async(src)
2385        .await
2386        .map_err(|err| data_error(format!("failed to stat '{}': {err}", src.display())))?;
2387    if !metadata.is_dir() {
2388        return Err(data_error(format!(
2389            "expected dataset directory at '{}'",
2390            src.display()
2391        )));
2392    }
2393    runmat_filesystem::create_dir_all_async(dst)
2394        .await
2395        .map_err(|err| data_error(format!("failed to create '{}': {err}", dst.display())))?;
2396    for entry in runmat_filesystem::read_dir_async(src)
2397        .await
2398        .map_err(|err| data_error(format!("failed to read '{}': {err}", src.display())))?
2399    {
2400        let entry_src = entry.path().to_path_buf();
2401        let entry_dst = dst.join(entry.file_name());
2402        if entry.is_dir() {
2403            copy_dir_recursive(&entry_src, &entry_dst).await?;
2404            continue;
2405        }
2406        copy_file(&entry_src, &entry_dst).await?;
2407    }
2408    Ok(())
2409}
2410
2411fn parse_shape_from_value(value: &Value) -> BuiltinResult<Vec<usize>> {
2412    match value {
2413        Value::Tensor(t) => {
2414            let mut out = Vec::with_capacity(t.data.len());
2415            for v in &t.data {
2416                if !v.is_finite() || *v < 0.0 {
2417                    return Err(data_error(
2418                        "shape dimensions must be non-negative finite numbers",
2419                    ));
2420                }
2421                out.push(*v as usize);
2422            }
2423            Ok(out)
2424        }
2425        Value::Num(v) => {
2426            if !v.is_finite() || *v < 0.0 {
2427                return Err(data_error(
2428                    "shape dimensions must be non-negative finite numbers",
2429                ));
2430            }
2431            Ok(vec![*v as usize])
2432        }
2433        Value::Int(v) => {
2434            let n = v.to_i64();
2435            if n < 0 {
2436                return Err(data_error("shape dimensions must be non-negative"));
2437            }
2438            Ok(vec![n as usize])
2439        }
2440        _ => Err(data_error("shape must be a numeric vector")),
2441    }
2442}
2443
2444async fn write_array_full_async(
2445    dataset_path: &str,
2446    array_name: &str,
2447    slice_spec: Option<&Value>,
2448    value: &Value,
2449) -> BuiltinResult<()> {
2450    let root = dataset_root(dataset_path);
2451    let mut manifest = read_manifest_async(&root).await?;
2452    apply_write_to_manifest_async(&root, &mut manifest, array_name, slice_spec, value).await?;
2453    manifest.updated_at = now_rfc3339();
2454    manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2455    write_manifest_async(&root, &manifest).await
2456}
2457
2458async fn apply_write_to_manifest_async(
2459    root: &std::path::Path,
2460    manifest: &mut DataManifest,
2461    array_name: &str,
2462    slice_spec: Option<&Value>,
2463    value: &Value,
2464) -> BuiltinResult<()> {
2465    let meta: DataArrayMeta = manifest
2466        .arrays
2467        .get(array_name)
2468        .cloned()
2469        .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2470
2471    if let Some(slice_spec) = slice_spec {
2472        if apply_slice_write_chunked_async(root, manifest, array_name, &meta, slice_spec, value)
2473            .await?
2474        {
2475            return Ok(());
2476        }
2477    }
2478
2479    let payload = read_array_payload_async(root, &meta).await?;
2480    let next_payload = if let Some(slice_spec) = slice_spec {
2481        write_slice_payload(&payload, slice_spec, value)?
2482    } else {
2483        DataArrayPayload::from_value(payload.dtype.clone(), value)?
2484    };
2485
2486    let (payload_path, chunk_index_path) =
2487        write_array_payload_async(root, array_name, &next_payload, &meta.chunk_shape).await?;
2488    if let Some(updated) = manifest.arrays.get_mut(array_name) {
2489        updated.shape = next_payload.shape.clone();
2490        updated.data_path = make_rel_data_path(root, &payload_path)?;
2491        updated.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2492    }
2493    Ok(())
2494}
2495
2496async fn apply_slice_write_chunked_async(
2497    root: &std::path::Path,
2498    manifest: &mut DataManifest,
2499    array_name: &str,
2500    meta: &DataArrayMeta,
2501    slice_spec: &Value,
2502    value: &Value,
2503) -> BuiltinResult<bool> {
2504    let Some(index_rel_path) = &meta.chunk_index_path else {
2505        return Ok(false);
2506    };
2507    let index_path = root.join(index_rel_path);
2508    if runmat_filesystem::metadata_async(&index_path)
2509        .await
2510        .is_err()
2511    {
2512        return Ok(false);
2513    }
2514    let ranges = parse_slice_spec(slice_spec, &meta.shape)?;
2515    let rhs_shape: Vec<usize> = ranges
2516        .iter()
2517        .map(|r| r.end.saturating_sub(r.start))
2518        .collect();
2519    let rhs = DataArrayPayload::from_value(meta.dtype.clone(), value)?;
2520    if rhs.shape != rhs_shape {
2521        return Err(data_error(format!(
2522            "SHAPE_MISMATCH: rhs shape {:?} must match target slice shape {:?}",
2523            rhs.shape, rhs_shape
2524        )));
2525    }
2526
2527    let index_bytes = runmat_filesystem::read_async(&index_path)
2528        .await
2529        .map_err(|err| {
2530            data_error(format!(
2531                "failed to read chunk index '{}': {err}",
2532                index_path.display()
2533            ))
2534        })?;
2535    let mut chunk_index: DataChunkIndex = serde_json::from_slice(&index_bytes).map_err(|err| {
2536        data_error(format!(
2537            "failed to parse chunk index '{}': {err}",
2538            index_path.display()
2539        ))
2540    })?;
2541
2542    let mut pos_by_key = HashMap::new();
2543    for (idx, entry) in chunk_index.chunks.iter().enumerate() {
2544        pos_by_key.insert(entry.key.clone(), idx);
2545    }
2546    let touched = touched_chunk_coords(&ranges, &meta.chunk_shape, &meta.shape);
2547    let mut upload_batch = Vec::<(DataChunkDescriptor, Vec<u8>)>::new();
2548
2549    for coords in touched {
2550        let key = chunk_key(&coords);
2551        let chunk_start = chunk_start_for_coords(&coords, &meta.chunk_shape);
2552        let chunk_extent = chunk_extent_for_start(&chunk_start, &meta.chunk_shape, &meta.shape);
2553        let intersection = chunk_intersection(&ranges, &chunk_start, &chunk_extent);
2554        if intersection.is_empty() {
2555            continue;
2556        }
2557
2558        let (entry_index, existed, mut entry, mut chunk_payload) = load_or_init_chunk(
2559            root,
2560            array_name,
2561            &key,
2562            &coords,
2563            &chunk_extent,
2564            &meta.dtype,
2565            &pos_by_key,
2566            &chunk_index,
2567        )
2568        .await?;
2569
2570        let mut local = vec![0usize; intersection.len()];
2571        let intersection_shape: Vec<usize> = intersection
2572            .iter()
2573            .map(|r| r.end.saturating_sub(r.start))
2574            .collect();
2575        loop {
2576            let mut global = Vec::with_capacity(intersection.len());
2577            for dim in 0..intersection.len() {
2578                global.push(intersection[dim].start + local[dim]);
2579            }
2580            let rhs_index: Vec<usize> = global
2581                .iter()
2582                .enumerate()
2583                .map(|(dim, g)| g.saturating_sub(ranges[dim].start))
2584                .collect();
2585            let chunk_index_local: Vec<usize> = global
2586                .iter()
2587                .enumerate()
2588                .map(|(dim, g)| g.saturating_sub(chunk_start[dim]))
2589                .collect();
2590            let rhs_linear = linear_index_column_major(&rhs_index, &rhs_shape)?;
2591            let chunk_linear = linear_index_column_major(&chunk_index_local, &chunk_extent)?;
2592            chunk_payload
2593                .values
2594                .set(chunk_linear, rhs.values.get(rhs_linear)?)?;
2595            if !advance_index(&mut local, &intersection_shape) {
2596                break;
2597            }
2598        }
2599
2600        let chunk_bytes = serde_json::to_vec(&chunk_payload)
2601            .map_err(|err| data_error(format!("failed to encode chunk payload: {err}")))?;
2602        let chunk_path = root.join(&entry.data_path);
2603        runmat_filesystem::write_async(&chunk_path, &chunk_bytes)
2604            .await
2605            .map_err(|err| {
2606                data_error(format!(
2607                    "failed to write chunk payload '{}': {err}",
2608                    chunk_path.display()
2609                ))
2610            })?;
2611
2612        entry.coords = coords.clone();
2613        entry.shape = chunk_extent.clone();
2614        entry.bytes_raw = chunk_bytes.len() as u64;
2615        entry.bytes_stored = chunk_bytes.len() as u64;
2616        entry.hash = sha256_hex(&chunk_bytes);
2617        if existed {
2618            chunk_index.chunks[entry_index] = entry.clone();
2619        } else {
2620            chunk_index.chunks.push(entry.clone());
2621            pos_by_key.insert(key.clone(), chunk_index.chunks.len() - 1);
2622        }
2623        upload_batch.push((
2624            DataChunkDescriptor {
2625                key: key.clone(),
2626                object_id: entry.object_id.clone(),
2627                hash: entry.hash.clone(),
2628                bytes_raw: entry.bytes_raw,
2629                bytes_stored: entry.bytes_stored,
2630            },
2631            chunk_bytes,
2632        ));
2633    }
2634
2635    maybe_upload_chunk_batch_async(root, array_name, upload_batch).await?;
2636    tracing::info!(
2637        target: "runmat.data",
2638        dataset = %root.display(),
2639        array = array_name,
2640        touched_chunks = chunk_index.chunks.len(),
2641        "chunked slice write committed"
2642    );
2643    let index_write = serde_json::to_vec(&chunk_index)
2644        .map_err(|err| data_error(format!("failed to encode chunk index json: {err}")))?;
2645    runmat_filesystem::write_async(&index_path, &index_write)
2646        .await
2647        .map_err(|err| {
2648            data_error(format!(
2649                "failed to write chunk index '{}': {err}",
2650                index_path.display()
2651            ))
2652        })?;
2653
2654    if let Some(updated) = manifest.arrays.get_mut(array_name) {
2655        updated.shape = meta.shape.clone();
2656        updated.chunk_index_path = Some(index_rel_path.clone());
2657    }
2658    Ok(true)
2659}
2660
2661#[derive(Clone, Copy, Debug)]
2662struct DimRange {
2663    start: usize,
2664    end: usize,
2665}
2666
2667fn read_slice_payload(
2668    payload: &DataArrayPayload,
2669    slice_spec: &Value,
2670) -> BuiltinResult<DataArrayPayload> {
2671    let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
2672    let out_shape: Vec<usize> = ranges
2673        .iter()
2674        .map(|r| r.end.saturating_sub(r.start))
2675        .collect();
2676    let mut out_values = crate::data::DataArrayValues::zeros(&payload.dtype, 0);
2677    let mut out_index = vec![0usize; out_shape.len()];
2678    loop {
2679        let source_index: Vec<usize> = out_index
2680            .iter()
2681            .enumerate()
2682            .map(|(dim, idx)| ranges[dim].start + *idx)
2683            .collect();
2684        let linear = linear_index_column_major(&source_index, &payload.shape)?;
2685        out_values.push(payload.values.get(linear)?)?;
2686
2687        if !advance_index(&mut out_index, &out_shape) {
2688            break;
2689        }
2690    }
2691    Ok(DataArrayPayload {
2692        dtype: payload.dtype.clone(),
2693        shape: out_shape,
2694        values: out_values,
2695    })
2696}
2697
2698fn write_slice_payload(
2699    payload: &DataArrayPayload,
2700    slice_spec: &Value,
2701    rhs: &Value,
2702) -> BuiltinResult<DataArrayPayload> {
2703    let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
2704    let target_shape: Vec<usize> = ranges
2705        .iter()
2706        .map(|r| r.end.saturating_sub(r.start))
2707        .collect();
2708    let rhs = DataArrayPayload::from_value(payload.dtype.clone(), rhs)?;
2709    if rhs.shape != target_shape {
2710        return Err(data_error(format!(
2711            "SHAPE_MISMATCH: rhs shape {:?} must match target slice shape {:?}",
2712            rhs.shape, target_shape
2713        )));
2714    }
2715
2716    let mut next = payload.values.clone();
2717    let mut rhs_index = vec![0usize; target_shape.len()];
2718    let mut rhs_linear = 0usize;
2719    loop {
2720        let target_index: Vec<usize> = rhs_index
2721            .iter()
2722            .enumerate()
2723            .map(|(dim, idx)| ranges[dim].start + *idx)
2724            .collect();
2725        let target_linear = linear_index_column_major(&target_index, &payload.shape)?;
2726        next.set(target_linear, rhs.values.get(rhs_linear)?)?;
2727        rhs_linear += 1;
2728
2729        if !advance_index(&mut rhs_index, &target_shape) {
2730            break;
2731        }
2732    }
2733
2734    Ok(DataArrayPayload {
2735        dtype: payload.dtype.clone(),
2736        shape: payload.shape.clone(),
2737        values: next,
2738    })
2739}
2740
2741fn parse_slice_spec(slice_spec: &Value, shape: &[usize]) -> BuiltinResult<Vec<DimRange>> {
2742    match slice_spec {
2743        Value::Cell(cell) => {
2744            if cell.data.is_empty() {
2745                return Err(data_error("INVALID_SLICE: empty slice specification"));
2746            }
2747            let mut ranges = Vec::with_capacity(shape.len());
2748            for (dim, extent) in shape.iter().enumerate() {
2749                if let Some(item) = cell.data.get(dim) {
2750                    ranges.push(parse_dim_range(item, *extent)?);
2751                } else {
2752                    ranges.push(DimRange {
2753                        start: 0,
2754                        end: *extent,
2755                    });
2756                }
2757            }
2758            Ok(ranges)
2759        }
2760        Value::String(s) if s == ":" => Ok(shape
2761            .iter()
2762            .map(|extent| DimRange {
2763                start: 0,
2764                end: *extent,
2765            })
2766            .collect()),
2767        _ => Err(data_error(
2768            "INVALID_SLICE: slice must be a cell spec like {1:10, :} or ':'",
2769        )),
2770    }
2771}
2772
2773fn parse_dim_range(value: &Value, extent: usize) -> BuiltinResult<DimRange> {
2774    if extent == 0 {
2775        return Ok(DimRange { start: 0, end: 0 });
2776    }
2777    match value {
2778        Value::String(s) if s == ":" => Ok(DimRange {
2779            start: 0,
2780            end: extent,
2781        }),
2782        Value::Num(n) => {
2783            let idx = (*n as isize) - 1;
2784            if idx < 0 || idx as usize >= extent {
2785                return Err(data_error("INVALID_SLICE: index out of bounds"));
2786            }
2787            Ok(DimRange {
2788                start: idx as usize,
2789                end: idx as usize + 1,
2790            })
2791        }
2792        Value::Int(i) => {
2793            let idx = i.to_i64() - 1;
2794            if idx < 0 || idx as usize >= extent {
2795                return Err(data_error("INVALID_SLICE: index out of bounds"));
2796            }
2797            Ok(DimRange {
2798                start: idx as usize,
2799                end: idx as usize + 1,
2800            })
2801        }
2802        Value::Tensor(t) if t.data.len() == 2 => {
2803            let start = (t.data[0] as isize) - 1;
2804            let end_inclusive = (t.data[1] as isize) - 1;
2805            if start < 0 || end_inclusive < start || end_inclusive as usize >= extent {
2806                return Err(data_error("INVALID_SLICE: range out of bounds"));
2807            }
2808            Ok(DimRange {
2809                start: start as usize,
2810                end: end_inclusive as usize + 1,
2811            })
2812        }
2813        _ => Err(data_error(
2814            "INVALID_SLICE: dimension must be ':', scalar index, or [start end] range",
2815        )),
2816    }
2817}
2818
2819fn linear_index_column_major(index: &[usize], shape: &[usize]) -> BuiltinResult<usize> {
2820    if index.len() != shape.len() {
2821        return Err(data_error("INVALID_SLICE: rank mismatch"));
2822    }
2823    let mut stride = 1usize;
2824    let mut linear = 0usize;
2825    for (idx, extent) in index.iter().zip(shape.iter()) {
2826        if *idx >= *extent {
2827            return Err(data_error("INVALID_SLICE: index out of bounds"));
2828        }
2829        linear += idx * stride;
2830        stride = stride.saturating_mul(*extent);
2831    }
2832    Ok(linear)
2833}
2834
2835fn advance_index(index: &mut [usize], shape: &[usize]) -> bool {
2836    if shape.is_empty() {
2837        return false;
2838    }
2839    for dim in 0..shape.len() {
2840        index[dim] += 1;
2841        if index[dim] < shape[dim] {
2842            return true;
2843        }
2844        index[dim] = 0;
2845    }
2846    false
2847}
2848
2849fn chunk_key(coords: &[usize]) -> String {
2850    coords
2851        .iter()
2852        .map(|v| v.to_string())
2853        .collect::<Vec<_>>()
2854        .join(".")
2855}
2856
2857fn chunk_start_for_coords(coords: &[usize], chunk_shape: &[usize]) -> Vec<usize> {
2858    coords
2859        .iter()
2860        .enumerate()
2861        .map(|(dim, coord)| coord * chunk_shape.get(dim).copied().unwrap_or(1).max(1))
2862        .collect()
2863}
2864
2865fn chunk_extent_for_start(start: &[usize], chunk_shape: &[usize], shape: &[usize]) -> Vec<usize> {
2866    start
2867        .iter()
2868        .enumerate()
2869        .map(|(dim, start)| {
2870            let chunk = chunk_shape.get(dim).copied().unwrap_or(1).max(1);
2871            let end = (*start + chunk).min(shape[dim]);
2872            end.saturating_sub(*start)
2873        })
2874        .collect()
2875}
2876
2877fn chunk_intersection(
2878    ranges: &[DimRange],
2879    chunk_start: &[usize],
2880    chunk_extent: &[usize],
2881) -> Vec<DimRange> {
2882    let mut out = Vec::with_capacity(ranges.len());
2883    for dim in 0..ranges.len() {
2884        let c_start = chunk_start[dim];
2885        let c_end = c_start + chunk_extent[dim];
2886        let start = ranges[dim].start.max(c_start);
2887        let end = ranges[dim].end.min(c_end);
2888        if start >= end {
2889            return Vec::new();
2890        }
2891        out.push(DimRange { start, end });
2892    }
2893    out
2894}
2895
2896fn touched_chunk_coords(
2897    ranges: &[DimRange],
2898    chunk_shape: &[usize],
2899    shape: &[usize],
2900) -> Vec<Vec<usize>> {
2901    let mut span = Vec::with_capacity(ranges.len());
2902    let mut begin = Vec::with_capacity(ranges.len());
2903    for dim in 0..ranges.len() {
2904        if shape[dim] == 0 {
2905            return Vec::new();
2906        }
2907        let chunk = chunk_shape.get(dim).copied().unwrap_or(1).max(1);
2908        let first = ranges[dim].start / chunk;
2909        let last = (ranges[dim].end.saturating_sub(1)) / chunk;
2910        begin.push(first);
2911        span.push(last.saturating_sub(first) + 1);
2912    }
2913    let mut local = vec![0usize; span.len()];
2914    let mut out = Vec::new();
2915    loop {
2916        out.push(
2917            local
2918                .iter()
2919                .enumerate()
2920                .map(|(dim, v)| begin[dim] + *v)
2921                .collect::<Vec<_>>(),
2922        );
2923        if !advance_index(&mut local, &span) {
2924            break;
2925        }
2926    }
2927    out
2928}
2929
2930async fn maybe_upload_chunk_batch_async(
2931    root: &std::path::Path,
2932    array_name: &str,
2933    batch: Vec<(DataChunkDescriptor, Vec<u8>)>,
2934) -> BuiltinResult<()> {
2935    if batch.is_empty() {
2936        return Ok(());
2937    }
2938    let request = DataChunkUploadRequest {
2939        dataset_path: root.to_string_lossy().to_string(),
2940        array: array_name.to_string(),
2941        chunks: batch.iter().map(|(d, _)| d.clone()).collect(),
2942    };
2943    let targets = match runmat_filesystem::data_chunk_upload_targets_async(&request).await {
2944        Ok(targets) => targets,
2945        Err(err) if err.kind() == std::io::ErrorKind::Unsupported => return Ok(()),
2946        Err(err) => {
2947            return Err(data_error(format!(
2948                "failed to request data chunk upload targets: {err}"
2949            )))
2950        }
2951    };
2952    for (descriptor, bytes) in batch {
2953        let target = targets
2954            .iter()
2955            .find(|t| t.key == descriptor.key)
2956            .ok_or_else(|| {
2957                data_error(format!(
2958                    "missing upload target for chunk '{}'",
2959                    descriptor.key
2960                ))
2961            })?;
2962        runmat_filesystem::data_upload_chunk_async(target, &bytes)
2963            .await
2964            .map_err(|err| {
2965                data_error(format!(
2966                    "failed to upload chunk '{}': {err}",
2967                    descriptor.key
2968                ))
2969            })?;
2970        tracing::info!(
2971            target: "runmat.data",
2972            dataset = %root.display(),
2973            array = array_name,
2974            chunk_key = descriptor.key,
2975            bytes = bytes.len(),
2976            "chunk upload completed"
2977        );
2978    }
2979    Ok(())
2980}
2981
2982fn chunk_rel_path(array_name: &str, object_id: &str) -> String {
2983    format!("arrays/{array_name}/chunks/{object_id}.json")
2984}
2985
2986async fn load_or_init_chunk(
2987    root: &std::path::Path,
2988    array_name: &str,
2989    key: &str,
2990    coords: &[usize],
2991    chunk_extent: &[usize],
2992    dtype: &str,
2993    pos_by_key: &HashMap<String, usize>,
2994    chunk_index: &DataChunkIndex,
2995) -> BuiltinResult<(usize, bool, DataChunkIndexEntry, DataArrayPayload)> {
2996    if let Some(index) = pos_by_key.get(key).copied() {
2997        let entry = chunk_index
2998            .chunks
2999            .get(index)
3000            .cloned()
3001            .ok_or_else(|| data_error(format!("chunk index missing key '{key}'")))?;
3002        let bytes = runmat_filesystem::read_async(root.join(&entry.data_path))
3003            .await
3004            .map_err(|err| {
3005                data_error(format!(
3006                    "failed to read chunk payload '{}': {err}",
3007                    entry.data_path
3008                ))
3009            })?;
3010        let payload: DataArrayPayload = serde_json::from_slice::<DataArrayPayload>(&bytes)
3011            .map_err(|err| {
3012                data_error(format!(
3013                    "failed to parse chunk payload '{}': {err}",
3014                    entry.data_path
3015                ))
3016            })?
3017            .normalize_for_dtype(dtype)?;
3018        return Ok((index, true, entry, payload));
3019    }
3020
3021    let object_id = format!("obj_{}", key.replace('.', "_"));
3022    let entry = DataChunkIndexEntry {
3023        key: key.to_string(),
3024        object_id: object_id.clone(),
3025        hash: String::new(),
3026        bytes_raw: 0,
3027        bytes_stored: 0,
3028        coords: coords.to_vec(),
3029        shape: chunk_extent.to_vec(),
3030        data_path: chunk_rel_path(array_name, &object_id),
3031    };
3032    let payload = DataArrayPayload::zeros(dtype.to_string(), chunk_extent.to_vec());
3033    Ok((chunk_index.chunks.len(), false, entry, payload))
3034}
3035
3036fn attrs_to_struct(attrs: &BTreeMap<String, serde_json::Value>) -> Value {
3037    let mut out = StructValue::new();
3038    for (k, v) in attrs {
3039        out.fields.insert(k.clone(), json_to_value(v));
3040    }
3041    Value::Struct(out)
3042}
3043
3044fn value_to_json(value: &Value) -> serde_json::Value {
3045    match value {
3046        Value::String(s) => serde_json::Value::String(s.clone()),
3047        Value::CharArray(chars) => serde_json::Value::String(chars.data.iter().collect::<String>()),
3048        Value::Num(n) => serde_json::json!(n),
3049        Value::Int(i) => serde_json::json!(i.to_i64()),
3050        Value::Bool(b) => serde_json::json!(b),
3051        _ => serde_json::Value::String(format!("{value:?}")),
3052    }
3053}
3054
3055fn json_to_value(value: &serde_json::Value) -> Value {
3056    match value {
3057        serde_json::Value::Bool(b) => Value::Bool(*b),
3058        serde_json::Value::Number(n) => Value::Num(n.as_f64().unwrap_or_default()),
3059        serde_json::Value::String(s) => Value::String(s.clone()),
3060        serde_json::Value::Array(arr) => {
3061            let vals = arr.iter().map(json_to_value).collect::<Vec<_>>();
3062            crate::make_cell(vals.clone(), 1, vals.len())
3063                .unwrap_or_else(|_| Value::String("<invalid-array>".to_string()))
3064        }
3065        serde_json::Value::Object(map) => {
3066            let mut s = StructValue::new();
3067            for (k, v) in map {
3068                s.fields.insert(k.clone(), json_to_value(v));
3069            }
3070            Value::Struct(s)
3071        }
3072        serde_json::Value::Null => Value::String("".to_string()),
3073    }
3074}
3075
3076#[cfg(all(test, not(target_arch = "wasm32")))]
3077mod tests {
3078    use super::*;
3079    use crate::dispatcher::call_builtin;
3080    use async_trait::async_trait;
3081    use axum::extract::{Query, State};
3082    use axum::http::{HeaderMap, StatusCode};
3083    use axum::routing::{post, put};
3084    use axum::{Json, Router};
3085    use runmat_builtins::{CellArray, IntegerStorage};
3086    use runmat_filesystem::data_contract::{
3087        DataChunkUploadRequest, DataChunkUploadTarget, DataManifestDescriptor, DataManifestRequest,
3088    };
3089    use runmat_filesystem::{
3090        DirEntry, FileHandle, FsMetadata, FsProvider, NativeFsProvider, OpenFlags,
3091    };
3092    use serde::Deserialize;
3093    use std::path::Path;
3094    use std::sync::{Arc, Mutex, MutexGuard};
3095    use tokio::runtime::Runtime;
3096    use tokio::sync::oneshot;
3097
3098    fn serial_test_guard() -> MutexGuard<'static, ()> {
3099        runmat_filesystem::provider_override_lock()
3100    }
3101
3102    fn native_provider_guard() -> runmat_filesystem::ProviderGuard {
3103        runmat_filesystem::replace_provider(Arc::new(NativeFsProvider))
3104    }
3105
3106    #[test]
3107    fn io_data_descriptors_cover_constructor_and_transaction_surface() {
3108        let data_labels: Vec<&str> = DATA_CREATE_DESCRIPTOR
3109            .signatures
3110            .iter()
3111            .map(|sig| sig.label)
3112            .collect();
3113        assert!(data_labels.contains(&"ds = data.create(path, schema, Name, Value, ...)"));
3114
3115        let read_labels: Vec<&str> = DATAARRAY_READ_DESCRIPTOR
3116            .signatures
3117            .iter()
3118            .map(|sig| sig.label)
3119            .collect();
3120        assert!(read_labels.contains(&"X = DataArray.read(arr, sliceSpec)"));
3121
3122        let write_labels: Vec<&str> = DATAARRAY_WRITE_DESCRIPTOR
3123            .signatures
3124            .iter()
3125            .map(|sig| sig.label)
3126            .collect();
3127        assert!(write_labels.contains(&"tf = DataArray.write(arr, values)"));
3128        assert!(write_labels.contains(&"tf = DataArray.write(arr, sliceSpec, values)"));
3129
3130        let tx_labels: Vec<&str> = DATATX_COMMIT_DESCRIPTOR
3131            .signatures
3132            .iter()
3133            .map(|sig| sig.label)
3134            .collect();
3135        assert!(tx_labels.contains(&"tf = DataTransaction.commit(tx, Name, Value, ...)"));
3136    }
3137
3138    #[derive(Default)]
3139    struct CountingDataUploadProvider {
3140        inner: NativeFsProvider,
3141        uploaded_keys: Arc<Mutex<Vec<String>>>,
3142    }
3143
3144    struct HttpDataUploadProvider {
3145        inner: NativeFsProvider,
3146        base_url: String,
3147        client: reqwest::blocking::Client,
3148    }
3149
3150    impl HttpDataUploadProvider {
3151        fn new(base_url: String) -> Self {
3152            Self {
3153                inner: NativeFsProvider,
3154                base_url,
3155                client: reqwest::blocking::Client::new(),
3156            }
3157        }
3158    }
3159
3160    #[async_trait(?Send)]
3161    impl FsProvider for HttpDataUploadProvider {
3162        fn open(&self, path: &Path, flags: &OpenFlags) -> std::io::Result<Box<dyn FileHandle>> {
3163            self.inner.open(path, flags)
3164        }
3165
3166        async fn read(&self, path: &Path) -> std::io::Result<Vec<u8>> {
3167            self.inner.read(path).await
3168        }
3169
3170        async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
3171            self.inner.write(path, data).await
3172        }
3173
3174        async fn remove_file(&self, path: &Path) -> std::io::Result<()> {
3175            self.inner.remove_file(path).await
3176        }
3177
3178        async fn metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
3179            self.inner.metadata(path).await
3180        }
3181
3182        async fn symlink_metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
3183            self.inner.symlink_metadata(path).await
3184        }
3185
3186        async fn read_dir(&self, path: &Path) -> std::io::Result<Vec<DirEntry>> {
3187            self.inner.read_dir(path).await
3188        }
3189
3190        async fn canonicalize(&self, path: &Path) -> std::io::Result<std::path::PathBuf> {
3191            self.inner.canonicalize(path).await
3192        }
3193
3194        async fn create_dir(&self, path: &Path) -> std::io::Result<()> {
3195            self.inner.create_dir(path).await
3196        }
3197
3198        async fn create_dir_all(&self, path: &Path) -> std::io::Result<()> {
3199            self.inner.create_dir_all(path).await
3200        }
3201
3202        async fn remove_dir(&self, path: &Path) -> std::io::Result<()> {
3203            self.inner.remove_dir(path).await
3204        }
3205
3206        async fn remove_dir_all(&self, path: &Path) -> std::io::Result<()> {
3207            self.inner.remove_dir_all(path).await
3208        }
3209
3210        async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
3211            self.inner.rename(from, to).await
3212        }
3213
3214        async fn set_readonly(&self, path: &Path, readonly: bool) -> std::io::Result<()> {
3215            self.inner.set_readonly(path, readonly).await
3216        }
3217
3218        async fn data_manifest_descriptor(
3219            &self,
3220            request: &DataManifestRequest,
3221        ) -> std::io::Result<DataManifestDescriptor> {
3222            self.inner.data_manifest_descriptor(request).await
3223        }
3224
3225        async fn data_chunk_upload_targets(
3226            &self,
3227            request: &DataChunkUploadRequest,
3228        ) -> std::io::Result<Vec<DataChunkUploadTarget>> {
3229            #[derive(Deserialize)]
3230            struct UploadTargetsResponse {
3231                targets: Vec<DataChunkUploadTarget>,
3232            }
3233            let url = format!("{}/data/chunks/upload-targets", self.base_url);
3234            let response = self
3235                .client
3236                .post(url)
3237                .json(request)
3238                .send()
3239                .map_err(|err| std::io::Error::other(err.to_string()))?;
3240            if !response.status().is_success() {
3241                return Err(std::io::Error::other(format!(
3242                    "upload targets request failed: {}",
3243                    response.status()
3244                )));
3245            }
3246            let parsed: UploadTargetsResponse = response
3247                .json()
3248                .map_err(|err| std::io::Error::other(err.to_string()))?;
3249            Ok(parsed.targets)
3250        }
3251
3252        async fn data_upload_chunk(
3253            &self,
3254            target: &DataChunkUploadTarget,
3255            data: &[u8],
3256        ) -> std::io::Result<()> {
3257            let upload_url = if let Some(key) = target.upload_url.strip_prefix("upload://") {
3258                format!("{}/upload?key={}", self.base_url, key)
3259            } else {
3260                target.upload_url.clone()
3261            };
3262            let method = reqwest::Method::from_bytes(target.method.as_bytes())
3263                .map_err(|err| std::io::Error::other(err.to_string()))?;
3264            let mut request = self.client.request(method, &upload_url);
3265            for (k, v) in &target.headers {
3266                request = request.header(k, v);
3267            }
3268            let response = request
3269                .body(data.to_vec())
3270                .send()
3271                .map_err(|err| std::io::Error::other(err.to_string()))?;
3272            if !response.status().is_success() {
3273                return Err(std::io::Error::other(format!(
3274                    "chunk upload failed: {}",
3275                    response.status()
3276                )));
3277            }
3278            Ok(())
3279        }
3280    }
3281
3282    #[derive(Clone, Default)]
3283    struct UploadHarness {
3284        uploads: Arc<Mutex<Vec<String>>>,
3285    }
3286
3287    #[derive(Deserialize)]
3288    struct UploadChunkQuery {
3289        key: String,
3290    }
3291
3292    async fn upload_targets_handler(
3293        Json(req): Json<DataChunkUploadRequest>,
3294    ) -> Result<Json<serde_json::Value>, StatusCode> {
3295        let targets = req
3296            .chunks
3297            .iter()
3298            .map(|chunk| {
3299                serde_json::json!({
3300                    "key": chunk.key,
3301                    "method": "PUT",
3302                    "upload_url": format!("upload://{}", chunk.key),
3303                    "headers": {
3304                        "x-runmat-hash": chunk.hash,
3305                    }
3306                })
3307            })
3308            .collect::<Vec<_>>();
3309        Ok(Json(serde_json::json!({ "targets": targets })))
3310    }
3311
3312    async fn upload_handler(
3313        State(harness): State<UploadHarness>,
3314        Query(query): Query<UploadChunkQuery>,
3315        headers: HeaderMap,
3316        body: axum::body::Bytes,
3317    ) -> Result<(), StatusCode> {
3318        if body.is_empty() {
3319            return Err(StatusCode::BAD_REQUEST);
3320        }
3321        if headers.get("x-runmat-hash").is_none() {
3322            return Err(StatusCode::BAD_REQUEST);
3323        }
3324        let mut guard = harness.uploads.lock().expect("uploads lock poisoned");
3325        guard.push(query.key);
3326        Ok(())
3327    }
3328
3329    fn spawn_upload_server() -> (
3330        String,
3331        Arc<Mutex<Vec<String>>>,
3332        Runtime,
3333        oneshot::Sender<()>,
3334    ) {
3335        let harness = UploadHarness::default();
3336        let uploads = Arc::clone(&harness.uploads);
3337        let runtime = Runtime::new().expect("tokio runtime");
3338        let (addr, shutdown_tx) = runtime.block_on(async move {
3339            let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
3340                .await
3341                .expect("bind upload server");
3342            let addr = listener.local_addr().expect("local addr");
3343            let app = Router::new()
3344                .route("/data/chunks/upload-targets", post(upload_targets_handler))
3345                .route("/upload", put(upload_handler))
3346                .with_state(harness);
3347            let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
3348            let server = axum::serve(listener, app).with_graceful_shutdown(async {
3349                let _ = shutdown_rx.await;
3350            });
3351            tokio::spawn(async move {
3352                let _ = server.await;
3353            });
3354            (addr, shutdown_tx)
3355        });
3356        (format!("http://{}", addr), uploads, runtime, shutdown_tx)
3357    }
3358
3359    impl CountingDataUploadProvider {
3360        fn uploaded_keys(&self) -> Arc<Mutex<Vec<String>>> {
3361            Arc::clone(&self.uploaded_keys)
3362        }
3363    }
3364
3365    #[async_trait(?Send)]
3366    impl FsProvider for CountingDataUploadProvider {
3367        fn open(&self, path: &Path, flags: &OpenFlags) -> std::io::Result<Box<dyn FileHandle>> {
3368            self.inner.open(path, flags)
3369        }
3370
3371        async fn read(&self, path: &Path) -> std::io::Result<Vec<u8>> {
3372            self.inner.read(path).await
3373        }
3374
3375        async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
3376            self.inner.write(path, data).await
3377        }
3378
3379        async fn remove_file(&self, path: &Path) -> std::io::Result<()> {
3380            self.inner.remove_file(path).await
3381        }
3382
3383        async fn metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
3384            self.inner.metadata(path).await
3385        }
3386
3387        async fn symlink_metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
3388            self.inner.symlink_metadata(path).await
3389        }
3390
3391        async fn read_dir(&self, path: &Path) -> std::io::Result<Vec<DirEntry>> {
3392            self.inner.read_dir(path).await
3393        }
3394
3395        async fn canonicalize(&self, path: &Path) -> std::io::Result<std::path::PathBuf> {
3396            self.inner.canonicalize(path).await
3397        }
3398
3399        async fn create_dir(&self, path: &Path) -> std::io::Result<()> {
3400            self.inner.create_dir(path).await
3401        }
3402
3403        async fn create_dir_all(&self, path: &Path) -> std::io::Result<()> {
3404            self.inner.create_dir_all(path).await
3405        }
3406
3407        async fn remove_dir(&self, path: &Path) -> std::io::Result<()> {
3408            self.inner.remove_dir(path).await
3409        }
3410
3411        async fn remove_dir_all(&self, path: &Path) -> std::io::Result<()> {
3412            self.inner.remove_dir_all(path).await
3413        }
3414
3415        async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
3416            self.inner.rename(from, to).await
3417        }
3418
3419        async fn set_readonly(&self, path: &Path, readonly: bool) -> std::io::Result<()> {
3420            self.inner.set_readonly(path, readonly).await
3421        }
3422
3423        async fn data_manifest_descriptor(
3424            &self,
3425            request: &DataManifestRequest,
3426        ) -> std::io::Result<DataManifestDescriptor> {
3427            self.inner.data_manifest_descriptor(request).await
3428        }
3429
3430        async fn data_chunk_upload_targets(
3431            &self,
3432            request: &DataChunkUploadRequest,
3433        ) -> std::io::Result<Vec<DataChunkUploadTarget>> {
3434            Ok(request
3435                .chunks
3436                .iter()
3437                .map(|chunk| DataChunkUploadTarget {
3438                    key: chunk.key.clone(),
3439                    method: "PUT".to_string(),
3440                    upload_url: format!("count://{}", chunk.object_id),
3441                    headers: std::collections::HashMap::new(),
3442                })
3443                .collect())
3444        }
3445
3446        async fn data_upload_chunk(
3447            &self,
3448            target: &DataChunkUploadTarget,
3449            _data: &[u8],
3450        ) -> std::io::Result<()> {
3451            let mut guard = match self.uploaded_keys.lock() {
3452                Ok(guard) => guard,
3453                Err(poisoned) => poisoned.into_inner(),
3454            };
3455            guard.push(target.key.clone());
3456            Ok(())
3457        }
3458    }
3459
3460    #[test]
3461    fn create_open_write_read_dataset() {
3462        let _serial = serial_test_guard();
3463        let _provider_guard = native_provider_guard();
3464        let dir = tempfile::tempdir().expect("tempdir");
3465        let path = dir.path().join("sample.data").to_string_lossy().to_string();
3466
3467        let mut array_meta = StructValue::new();
3468        array_meta
3469            .fields
3470            .insert("dtype".to_string(), Value::String("f64".to_string()));
3471        array_meta.fields.insert(
3472            "shape".to_string(),
3473            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape tensor")),
3474        );
3475        let mut arrays = StructValue::new();
3476        arrays
3477            .fields
3478            .insert("temperature".to_string(), Value::Struct(array_meta));
3479        let mut schema = StructValue::new();
3480        schema
3481            .fields
3482            .insert("arrays".to_string(), Value::Struct(arrays));
3483
3484        let ds = call_builtin(
3485            "data.create",
3486            &[
3487                Value::String(path.clone()),
3488                Value::Struct(schema),
3489                Value::Cell(runmat_builtins::CellArray::new(vec![], 1, 0).expect("cell")),
3490            ],
3491        )
3492        .expect("create dataset");
3493
3494        let arr = call_builtin(
3495            "Dataset.array",
3496            &[ds, Value::String("temperature".to_string())],
3497        )
3498        .expect("dataset array");
3499        let write_tensor = Tensor::new(vec![1.0, 2.0, 3.0, 4.0], vec![2, 2]).expect("write tensor");
3500        call_builtin(
3501            "DataArray.write",
3502            &[arr.clone(), Value::Tensor(write_tensor)],
3503        )
3504        .expect("write array");
3505
3506        let read_back = call_builtin("DataArray.read", &[arr]).expect("read array");
3507        let Value::Tensor(t) = read_back else {
3508            panic!("expected tensor");
3509        };
3510        assert_eq!(t.shape, vec![2, 2]);
3511        assert_eq!(t.data, vec![1.0, 2.0, 3.0, 4.0]);
3512    }
3513
3514    #[test]
3515    fn write_and_read_slice_payload() {
3516        let _serial = serial_test_guard();
3517        let _provider_guard = native_provider_guard();
3518        let dir = tempfile::tempdir().expect("tempdir");
3519        let path = dir.path().join("slice.data").to_string_lossy().to_string();
3520
3521        let mut array_meta = StructValue::new();
3522        array_meta
3523            .fields
3524            .insert("dtype".to_string(), Value::String("f64".to_string()));
3525        array_meta.fields.insert(
3526            "shape".to_string(),
3527            Value::Tensor(Tensor::new(vec![3.0, 3.0], vec![1, 2]).expect("shape tensor")),
3528        );
3529        let mut arrays = StructValue::new();
3530        arrays
3531            .fields
3532            .insert("temperature".to_string(), Value::Struct(array_meta));
3533        let mut schema = StructValue::new();
3534        schema
3535            .fields
3536            .insert("arrays".to_string(), Value::Struct(arrays));
3537
3538        let ds = call_builtin(
3539            "data.create",
3540            &[
3541                Value::String(path.clone()),
3542                Value::Struct(schema),
3543                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3544            ],
3545        )
3546        .expect("create dataset");
3547        let arr = call_builtin(
3548            "Dataset.array",
3549            &[ds, Value::String("temperature".to_string())],
3550        )
3551        .expect("dataset array");
3552
3553        let slice = Value::Cell(
3554            CellArray::new(
3555                vec![
3556                    Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
3557                    Value::String(":".to_string()),
3558                ],
3559                1,
3560                2,
3561            )
3562            .expect("slice cell"),
3563        );
3564        let rhs = Value::Tensor(
3565            Tensor::new(vec![10.0, 11.0, 12.0, 13.0, 14.0, 15.0], vec![2, 3]).expect("rhs"),
3566        );
3567        call_builtin("DataArray.write", &[arr.clone(), slice.clone(), rhs]).expect("slice write");
3568
3569        let read_back = call_builtin("DataArray.read", &[arr.clone(), slice]).expect("slice read");
3570        let Value::Tensor(t) = read_back else {
3571            panic!("expected tensor");
3572        };
3573        assert_eq!(t.shape, vec![2, 3]);
3574        assert_eq!(t.data, vec![10.0, 11.0, 12.0, 13.0, 14.0, 15.0]);
3575    }
3576
3577    #[test]
3578    fn slice_write_updates_only_touched_chunks() {
3579        let _serial = serial_test_guard();
3580        let _provider_guard = native_provider_guard();
3581        let dir = tempfile::tempdir().expect("tempdir");
3582        let path = dir
3583            .path()
3584            .join("chunked.data")
3585            .to_string_lossy()
3586            .to_string();
3587
3588        let mut array_meta = StructValue::new();
3589        array_meta
3590            .fields
3591            .insert("dtype".to_string(), Value::String("f64".to_string()));
3592        array_meta.fields.insert(
3593            "shape".to_string(),
3594            Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
3595        );
3596        array_meta.fields.insert(
3597            "chunk".to_string(),
3598            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
3599        );
3600        let mut arrays = StructValue::new();
3601        arrays
3602            .fields
3603            .insert("temperature".to_string(), Value::Struct(array_meta));
3604        let mut schema = StructValue::new();
3605        schema
3606            .fields
3607            .insert("arrays".to_string(), Value::Struct(arrays));
3608
3609        let ds = call_builtin(
3610            "data.create",
3611            &[
3612                Value::String(path.clone()),
3613                Value::Struct(schema),
3614                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3615            ],
3616        )
3617        .expect("create dataset");
3618        let arr = call_builtin(
3619            "Dataset.array",
3620            &[ds, Value::String("temperature".to_string())],
3621        )
3622        .expect("dataset array");
3623
3624        let full = Value::Tensor(
3625            Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4]).expect("full tensor"),
3626        );
3627        call_builtin("DataArray.write", &[arr.clone(), full]).expect("initial write");
3628
3629        let root = std::path::PathBuf::from(&path);
3630        let untouched_path = root.join("arrays/temperature/chunks/obj_1_1.json");
3631        let touched_path = root.join("arrays/temperature/chunks/obj_0_0.json");
3632        let untouched_before =
3633            futures::executor::block_on(runmat_filesystem::read_async(&untouched_path))
3634                .expect("read untouched before");
3635        let touched_before =
3636            futures::executor::block_on(runmat_filesystem::read_async(&touched_path))
3637                .expect("read touched before");
3638
3639        let slice = Value::Cell(
3640            CellArray::new(
3641                vec![
3642                    Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
3643                    Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
3644                ],
3645                1,
3646                2,
3647            )
3648            .expect("slice cell"),
3649        );
3650        let rhs =
3651            Value::Tensor(Tensor::new(vec![99.0, 98.0, 97.0, 96.0], vec![2, 2]).expect("rhs"));
3652        call_builtin("DataArray.write", &[arr.clone(), slice, rhs]).expect("slice write");
3653
3654        let untouched_after =
3655            futures::executor::block_on(runmat_filesystem::read_async(&untouched_path))
3656                .expect("read untouched after");
3657        let touched_after =
3658            futures::executor::block_on(runmat_filesystem::read_async(&touched_path))
3659                .expect("read touched after");
3660        assert_eq!(untouched_before, untouched_after);
3661        assert_ne!(touched_before, touched_after);
3662    }
3663
3664    #[test]
3665    fn slice_write_uploads_only_touched_chunk_targets() {
3666        let _serial = serial_test_guard();
3667        let provider = Arc::new(CountingDataUploadProvider::default());
3668        let uploaded = provider.uploaded_keys();
3669        let _guard = runmat_filesystem::replace_provider(provider);
3670
3671        let dir = tempfile::tempdir().expect("tempdir");
3672        let path = dir
3673            .path()
3674            .join("remote-chunked.data")
3675            .to_string_lossy()
3676            .to_string();
3677
3678        let mut array_meta = StructValue::new();
3679        array_meta
3680            .fields
3681            .insert("dtype".to_string(), Value::String("f64".to_string()));
3682        array_meta.fields.insert(
3683            "shape".to_string(),
3684            Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
3685        );
3686        array_meta.fields.insert(
3687            "chunk".to_string(),
3688            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
3689        );
3690        let mut arrays = StructValue::new();
3691        arrays
3692            .fields
3693            .insert("temperature".to_string(), Value::Struct(array_meta));
3694        let mut schema = StructValue::new();
3695        schema
3696            .fields
3697            .insert("arrays".to_string(), Value::Struct(arrays));
3698
3699        let ds = call_builtin(
3700            "data.create",
3701            &[
3702                Value::String(path.clone()),
3703                Value::Struct(schema),
3704                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3705            ],
3706        )
3707        .expect("create dataset");
3708        let arr = call_builtin(
3709            "Dataset.array",
3710            &[ds, Value::String("temperature".to_string())],
3711        )
3712        .expect("dataset array");
3713
3714        call_builtin(
3715            "DataArray.write",
3716            &[
3717                arr.clone(),
3718                Value::Tensor(
3719                    Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
3720                        .expect("full tensor"),
3721                ),
3722            ],
3723        )
3724        .expect("initial write");
3725
3726        let manifest =
3727            futures::executor::block_on(crate::data::read_manifest_async(&dataset_root(&path)))
3728                .expect("manifest after initial write");
3729        let meta = manifest
3730            .arrays
3731            .get("temperature")
3732            .expect("temperature meta");
3733        let chunk_index_path =
3734            dataset_root(&path).join(meta.chunk_index_path.clone().expect("chunk index path"));
3735        assert!(
3736            futures::executor::block_on(runmat_filesystem::metadata_async(&chunk_index_path))
3737                .is_ok()
3738        );
3739
3740        {
3741            let mut keys = uploaded.lock().expect("uploaded keys lock");
3742            keys.clear();
3743        }
3744
3745        let slice = Value::Cell(
3746            CellArray::new(
3747                vec![
3748                    Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
3749                    Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
3750                ],
3751                1,
3752                2,
3753            )
3754            .expect("slice cell"),
3755        );
3756        let rhs = Value::Tensor(Tensor::new(vec![9.0, 8.0, 7.0, 6.0], vec![2, 2]).expect("rhs"));
3757        call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
3758
3759        let keys = uploaded.lock().expect("uploaded keys lock");
3760        assert_eq!(keys.as_slice(), ["0.0".to_string()].as_slice());
3761    }
3762
3763    #[test]
3764    fn slice_write_uploads_expected_cross_boundary_chunk_targets() {
3765        let _serial = serial_test_guard();
3766        let provider = Arc::new(CountingDataUploadProvider::default());
3767        let uploaded = provider.uploaded_keys();
3768        let _guard = runmat_filesystem::replace_provider(provider);
3769
3770        let dir = tempfile::tempdir().expect("tempdir");
3771        let path = dir
3772            .path()
3773            .join("remote-chunked-boundary.data")
3774            .to_string_lossy()
3775            .to_string();
3776
3777        let mut array_meta = StructValue::new();
3778        array_meta
3779            .fields
3780            .insert("dtype".to_string(), Value::String("f64".to_string()));
3781        array_meta.fields.insert(
3782            "shape".to_string(),
3783            Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
3784        );
3785        array_meta.fields.insert(
3786            "chunk".to_string(),
3787            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
3788        );
3789        let mut arrays = StructValue::new();
3790        arrays
3791            .fields
3792            .insert("temperature".to_string(), Value::Struct(array_meta));
3793        let mut schema = StructValue::new();
3794        schema
3795            .fields
3796            .insert("arrays".to_string(), Value::Struct(arrays));
3797
3798        let ds = call_builtin(
3799            "data.create",
3800            &[
3801                Value::String(path.clone()),
3802                Value::Struct(schema),
3803                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3804            ],
3805        )
3806        .expect("create dataset");
3807        let arr = call_builtin(
3808            "Dataset.array",
3809            &[ds, Value::String("temperature".to_string())],
3810        )
3811        .expect("dataset array");
3812
3813        call_builtin(
3814            "DataArray.write",
3815            &[
3816                arr.clone(),
3817                Value::Tensor(
3818                    Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
3819                        .expect("full tensor"),
3820                ),
3821            ],
3822        )
3823        .expect("initial write");
3824        {
3825            let mut keys = uploaded.lock().expect("uploaded keys lock");
3826            keys.clear();
3827        }
3828
3829        let slice = Value::Cell(
3830            CellArray::new(
3831                vec![
3832                    Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
3833                    Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
3834                ],
3835                1,
3836                2,
3837            )
3838            .expect("slice cell"),
3839        );
3840        let rhs =
3841            Value::Tensor(Tensor::new(vec![19.0, 18.0, 17.0, 16.0], vec![2, 2]).expect("rhs"));
3842        call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
3843
3844        let mut keys = uploaded.lock().expect("uploaded keys lock").clone();
3845        keys.sort();
3846        keys.dedup();
3847        assert_eq!(
3848            keys.as_slice(),
3849            [
3850                "0.0".to_string(),
3851                "0.1".to_string(),
3852                "1.0".to_string(),
3853                "1.1".to_string(),
3854            ]
3855            .as_slice()
3856        );
3857    }
3858
3859    #[test]
3860    fn slice_write_hits_http_server_data_endpoints_with_expected_keys() {
3861        let _serial = serial_test_guard();
3862        let (base_url, uploads, runtime, shutdown_tx) = spawn_upload_server();
3863        let provider = Arc::new(HttpDataUploadProvider::new(base_url));
3864        let _guard = runmat_filesystem::replace_provider(provider);
3865
3866        let dir = tempfile::tempdir().expect("tempdir");
3867        let path = dir
3868            .path()
3869            .join("http-endpoint.data")
3870            .to_string_lossy()
3871            .to_string();
3872
3873        let mut array_meta = StructValue::new();
3874        array_meta
3875            .fields
3876            .insert("dtype".to_string(), Value::String("f64".to_string()));
3877        array_meta.fields.insert(
3878            "shape".to_string(),
3879            Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
3880        );
3881        array_meta.fields.insert(
3882            "chunk".to_string(),
3883            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
3884        );
3885        let mut arrays = StructValue::new();
3886        arrays
3887            .fields
3888            .insert("temperature".to_string(), Value::Struct(array_meta));
3889        let mut schema = StructValue::new();
3890        schema
3891            .fields
3892            .insert("arrays".to_string(), Value::Struct(arrays));
3893
3894        let ds = call_builtin(
3895            "data.create",
3896            &[
3897                Value::String(path.clone()),
3898                Value::Struct(schema),
3899                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3900            ],
3901        )
3902        .expect("create dataset");
3903        let arr = call_builtin(
3904            "Dataset.array",
3905            &[ds, Value::String("temperature".to_string())],
3906        )
3907        .expect("dataset array");
3908
3909        call_builtin(
3910            "DataArray.write",
3911            &[
3912                arr.clone(),
3913                Value::Tensor(
3914                    Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
3915                        .expect("full tensor"),
3916                ),
3917            ],
3918        )
3919        .expect("initial write");
3920
3921        {
3922            let mut keys = uploads.lock().expect("uploads lock");
3923            keys.clear();
3924        }
3925
3926        let slice = Value::Cell(
3927            CellArray::new(
3928                vec![
3929                    Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
3930                    Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
3931                ],
3932                1,
3933                2,
3934            )
3935            .expect("slice cell"),
3936        );
3937        let rhs =
3938            Value::Tensor(Tensor::new(vec![19.0, 18.0, 17.0, 16.0], vec![2, 2]).expect("rhs"));
3939        call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
3940
3941        let mut keys = uploads.lock().expect("uploads lock").clone();
3942        keys.sort();
3943        keys.dedup();
3944        assert_eq!(
3945            keys.as_slice(),
3946            [
3947                "0.0".to_string(),
3948                "0.1".to_string(),
3949                "1.0".to_string(),
3950                "1.1".to_string(),
3951            ]
3952            .as_slice()
3953        );
3954
3955        let _ = shutdown_tx.send(());
3956        drop(runtime);
3957    }
3958
3959    #[test]
3960    fn tx_create_resize_fill_and_delete_array() {
3961        let _serial = serial_test_guard();
3962        let _provider = native_provider_guard();
3963        let dir = tempfile::tempdir().expect("tempdir");
3964        let path = dir.path().join("tx-ops.data").to_string_lossy().to_string();
3965
3966        let mut arrays = StructValue::new();
3967        let mut array_meta = StructValue::new();
3968        array_meta
3969            .fields
3970            .insert("dtype".to_string(), Value::String("f64".to_string()));
3971        array_meta.fields.insert(
3972            "shape".to_string(),
3973            Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("shape tensor")),
3974        );
3975        arrays
3976            .fields
3977            .insert("base".to_string(), Value::Struct(array_meta));
3978        let mut schema = StructValue::new();
3979        schema
3980            .fields
3981            .insert("arrays".to_string(), Value::Struct(arrays));
3982
3983        let ds = call_builtin(
3984            "data.create",
3985            &[
3986                Value::String(path.clone()),
3987                Value::Struct(schema),
3988                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
3989            ],
3990        )
3991        .expect("create dataset");
3992
3993        let tx = call_builtin("Dataset.begin", &[ds]).expect("begin tx");
3994        let mut new_meta = StructValue::new();
3995        new_meta
3996            .fields
3997            .insert("dtype".to_string(), Value::String("f64".to_string()));
3998        new_meta.fields.insert(
3999            "shape".to_string(),
4000            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape tensor")),
4001        );
4002        call_builtin(
4003            "DataTransaction.create_array",
4004            &[
4005                tx.clone(),
4006                Value::String("new_array".to_string()),
4007                Value::Struct(new_meta),
4008            ],
4009        )
4010        .expect("create array in tx");
4011        call_builtin(
4012            "DataTransaction.resize",
4013            &[
4014                tx.clone(),
4015                Value::String("new_array".to_string()),
4016                Value::Tensor(Tensor::new(vec![3.0, 1.0], vec![1, 2]).expect("shape tensor")),
4017            ],
4018        )
4019        .expect("resize array in tx");
4020        call_builtin(
4021            "DataTransaction.fill",
4022            &[
4023                tx.clone(),
4024                Value::String("new_array".to_string()),
4025                Value::Num(7.0),
4026            ],
4027        )
4028        .expect("fill array in tx");
4029        call_builtin(
4030            "DataTransaction.delete_array",
4031            &[tx.clone(), Value::String("base".to_string())],
4032        )
4033        .expect("delete array in tx");
4034        call_builtin("DataTransaction.commit", &[tx]).expect("commit tx");
4035
4036        let ds = call_builtin(
4037            "data.open",
4038            &[
4039                Value::String(path),
4040                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
4041            ],
4042        )
4043        .expect("open dataset");
4044        let has_base = call_builtin(
4045            "Dataset.has_array",
4046            &[ds.clone(), Value::String("base".to_string())],
4047        )
4048        .expect("has base");
4049        assert_eq!(has_base, Value::Bool(false));
4050        let arr = call_builtin(
4051            "Dataset.array",
4052            &[ds, Value::String("new_array".to_string())],
4053        )
4054        .expect("new array");
4055        let read_back = call_builtin("DataArray.read", &[arr]).expect("read array");
4056        let Value::Tensor(t) = read_back else {
4057            panic!("expected tensor");
4058        };
4059        assert_eq!(t.shape, vec![3, 1]);
4060        assert_eq!(t.data, vec![7.0, 7.0, 7.0]);
4061    }
4062
4063    #[test]
4064    fn data_arrays_preserve_every_integer_class_through_chunked_write_and_read() {
4065        let _serial = serial_test_guard();
4066        let _provider = native_provider_guard();
4067        let cases = vec![
4068            ("int8", IntegerStorage::I8(vec![i8::MIN, i8::MAX])),
4069            ("int16", IntegerStorage::I16(vec![i16::MIN, i16::MAX])),
4070            ("int32", IntegerStorage::I32(vec![i32::MIN, i32::MAX])),
4071            ("int64", IntegerStorage::I64(vec![i64::MIN, i64::MAX])),
4072            ("uint8", IntegerStorage::U8(vec![0, u8::MAX])),
4073            ("uint16", IntegerStorage::U16(vec![0, u16::MAX])),
4074            ("uint32", IntegerStorage::U32(vec![0, u32::MAX])),
4075            ("uint64", IntegerStorage::U64(vec![0, u64::MAX])),
4076        ];
4077
4078        for (dtype, storage) in cases {
4079            let dir = tempfile::tempdir().expect("tempdir");
4080            let path = dir
4081                .path()
4082                .join(format!("{dtype}.data"))
4083                .to_string_lossy()
4084                .to_string();
4085            let mut array_meta = StructValue::new();
4086            array_meta
4087                .fields
4088                .insert("dtype".to_string(), Value::String(dtype.to_string()));
4089            array_meta.fields.insert(
4090                "shape".to_string(),
4091                Value::Tensor(Tensor::new(vec![2.0, 1.0], vec![1, 2]).expect("shape")),
4092            );
4093            array_meta.fields.insert(
4094                "chunk".to_string(),
4095                Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("chunk")),
4096            );
4097            let mut arrays = StructValue::new();
4098            arrays
4099                .fields
4100                .insert("samples".to_string(), Value::Struct(array_meta));
4101            let mut schema = StructValue::new();
4102            schema
4103                .fields
4104                .insert("arrays".to_string(), Value::Struct(arrays));
4105
4106            let ds = call_builtin(
4107                "data.create",
4108                &[
4109                    Value::String(path),
4110                    Value::Struct(schema),
4111                    Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
4112                ],
4113            )
4114            .expect("create dataset");
4115            let arr = call_builtin("Dataset.array", &[ds, Value::String("samples".to_string())])
4116                .expect("array");
4117            let input = Tensor::new_integer(storage.clone(), vec![2, 1]).expect("integer tensor");
4118            call_builtin("DataArray.write", &[arr.clone(), Value::Tensor(input)])
4119                .expect("write integer array");
4120
4121            let Value::Tensor(read_back) = call_builtin("DataArray.read", &[arr]).expect("read")
4122            else {
4123                panic!("expected tensor");
4124            };
4125            assert_eq!(read_back.integer_storage(), Some(&storage), "{dtype}");
4126        }
4127    }
4128
4129    #[test]
4130    fn uint64_data_array_slice_fill_and_transaction_paths_remain_exact() {
4131        let _serial = serial_test_guard();
4132        let _provider = native_provider_guard();
4133        let dir = tempfile::tempdir().expect("tempdir");
4134        let path = dir.path().join("uint64.data").to_string_lossy().to_string();
4135        let mut array_meta = StructValue::new();
4136        array_meta
4137            .fields
4138            .insert("dtype".to_string(), Value::String("uint64".to_string()));
4139        array_meta.fields.insert(
4140            "shape".to_string(),
4141            Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape")),
4142        );
4143        array_meta.fields.insert(
4144            "chunk".to_string(),
4145            Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("chunk")),
4146        );
4147        let mut arrays = StructValue::new();
4148        arrays
4149            .fields
4150            .insert("samples".to_string(), Value::Struct(array_meta));
4151        let mut schema = StructValue::new();
4152        schema
4153            .fields
4154            .insert("arrays".to_string(), Value::Struct(arrays));
4155        let ds = call_builtin(
4156            "data.create",
4157            &[
4158                Value::String(path),
4159                Value::Struct(schema),
4160                Value::Cell(CellArray::new(vec![], 1, 0).expect("cell")),
4161            ],
4162        )
4163        .expect("create dataset");
4164        let arr = call_builtin(
4165            "Dataset.array",
4166            &[ds.clone(), Value::String("samples".to_string())],
4167        )
4168        .expect("array");
4169
4170        call_builtin(
4171            "DataArray.fill",
4172            &[
4173                arr.clone(),
4174                Value::Int(runmat_builtins::IntValue::U64(u64::MAX)),
4175            ],
4176        )
4177        .expect("fill");
4178        let slice = Value::Cell(
4179            CellArray::new(
4180                vec![
4181                    Value::Int(runmat_builtins::IntValue::I32(1)),
4182                    Value::String(":".to_string()),
4183                ],
4184                1,
4185                2,
4186            )
4187            .expect("slice"),
4188        );
4189        let replacement = Tensor::new_integer(
4190            IntegerStorage::U64(vec![1_u64 << 63, u64::MAX - 1]),
4191            vec![1, 2],
4192        )
4193        .expect("replacement");
4194        call_builtin(
4195            "DataArray.write",
4196            &[arr.clone(), slice, Value::Tensor(replacement)],
4197        )
4198        .expect("slice write");
4199
4200        let tx = call_builtin("Dataset.begin", &[ds]).expect("begin transaction");
4201        call_builtin(
4202            "DataTransaction.fill",
4203            &[
4204                tx.clone(),
4205                Value::String("samples".to_string()),
4206                Value::Int(runmat_builtins::IntValue::U64(1_u64 << 63)),
4207            ],
4208        )
4209        .expect("queue transaction fill");
4210        call_builtin("DataTransaction.commit", &[tx]).expect("commit transaction");
4211
4212        let Value::Tensor(read_back) = call_builtin("DataArray.read", &[arr]).expect("read") else {
4213            panic!("expected tensor");
4214        };
4215        assert_eq!(
4216            read_back.integer_storage(),
4217            Some(&IntegerStorage::U64(vec![1_u64 << 63; 4]))
4218        );
4219    }
4220}