1use std::collections::BTreeMap;
4use std::collections::HashMap;
5use std::path::PathBuf;
6
7use runmat_builtins::{
8 BuiltinCompletionPolicy, BuiltinDescriptor, BuiltinErrorDescriptor,
9 BuiltinIntegerAuditDescriptor, BuiltinIntegerAuditKind, BuiltinIntegerBackendRule,
10 BuiltinIntegerCapabilityDescriptor, BuiltinIntegerComputationDomain,
11 BuiltinIntegerInputAvailability, BuiltinIntegerInputCapability, BuiltinIntegerOutputClassRule,
12 BuiltinIntegerOverflowRule, BuiltinIntegerOverloadKind, BuiltinIntegerScalarDoubleRule,
13 BuiltinOutputMode, BuiltinParamArity, BuiltinParamDescriptor, BuiltinParamType,
14 BuiltinSignatureDescriptor,
15};
16use runmat_filesystem::data_contract::{DataChunkDescriptor, DataChunkUploadRequest};
17use runmat_macros::runtime_builtin;
18use runmat_value::{IntValue, NumericScalar, ObjectInstance, StructValue, Tensor, Value};
19
20use crate::builtins::common::json::int_value_to_json;
21use crate::builtins::common::spec::{
22 BroadcastSemantics, BuiltinFusionSpec, BuiltinGpuSpec, ConstantStrategy, GpuOpKind,
23 ReductionNaN, ResidencyPolicy, ShapeRequirements,
24};
25use crate::builtins::common::tensor as tensor_utils;
26use crate::data::{
27 array_object, data_error, dataset_object, dataset_root, ensure_manifest_sequence,
28 get_object_prop, manifest_path, manifest_version_token, now_rfc3339, parse_schema,
29 parse_string, read_array_payload_async, read_manifest_async, remove_tx, sha256_hex, start_tx,
30 transaction_object, validate_chunk_shape, with_tx, with_tx_mut, write_array_payload_async,
31 write_manifest_async, DataArrayMeta, DataArrayPayload, DataChunkIndex, DataChunkIndexEntry,
32 DataManifest, PendingCreateArray, PendingFill, PendingResize, PendingWrite, TxnStatus,
33};
34use crate::{make_cell, BuiltinResult};
35
36#[runmat_macros::register_gpu_spec(builtin_path = "crate::builtins::io::data")]
37pub const GPU_SPEC: BuiltinGpuSpec = BuiltinGpuSpec {
38 name: "data.*",
39 op_kind: GpuOpKind::Custom("io-data"),
40 supported_precisions: &[],
41 broadcast: BroadcastSemantics::None,
42 provider_hooks: &[],
43 constant_strategy: ConstantStrategy::InlineLiteral,
44 residency: ResidencyPolicy::GatherImmediately,
45 nan_mode: ReductionNaN::Include,
46 two_pass_threshold: None,
47 workgroup_size: None,
48 accepts_nan_mode: false,
49 notes: "Dataset operations are host I/O and metadata orchestration.",
50};
51
52#[runmat_macros::register_fusion_spec(builtin_path = "crate::builtins::io::data")]
53pub const FUSION_SPEC: BuiltinFusionSpec = BuiltinFusionSpec {
54 name: "data.*",
55 shape: ShapeRequirements::Any,
56 constant_strategy: ConstantStrategy::InlineLiteral,
57 elementwise: None,
58 reduction: None,
59 emits_nan: false,
60 notes: "Data builtins are side-effecting and not fusible.",
61};
62
63const DATA_ERROR_INVALID_ARGUMENT: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
64 code: "RM.DATA.INVALID_ARGUMENT",
65 identifier: Some("RunMat:data:InvalidArgument"),
66 when: "Arguments, receiver object, or option grammar are invalid for the requested data API.",
67 message: "data: invalid argument",
68};
69
70const DATA_ERROR_NOT_FOUND: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
71 code: "RM.DATA.NOT_FOUND",
72 identifier: Some("RunMat:data:NotFound"),
73 when: "Referenced dataset/array/transaction object is missing or cannot be resolved.",
74 message: "data: requested object not found",
75};
76
77const DATA_ERROR_INTERNAL: BuiltinErrorDescriptor = BuiltinErrorDescriptor {
78 code: "RM.DATA.INTERNAL",
79 identifier: Some("RunMat:data:Internal"),
80 when: "Filesystem/manifest/chunk processing fails unexpectedly during data API execution.",
81 message: "data: internal operation failed",
82};
83
84const DATA_DESCRIPTOR_ERRORS: [BuiltinErrorDescriptor; 3] = [
85 DATA_ERROR_INVALID_ARGUMENT,
86 DATA_ERROR_NOT_FOUND,
87 DATA_ERROR_INTERNAL,
88];
89
90const OUT_DATASET: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
91 name: "ds",
92 ty: BuiltinParamType::Any,
93 arity: BuiltinParamArity::Required,
94 default: None,
95 description: "Dataset handle object.",
96}];
97
98const OUT_ARRAY: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
99 name: "arr",
100 ty: BuiltinParamType::Any,
101 arity: BuiltinParamArity::Required,
102 default: None,
103 description: "DataArray handle object.",
104}];
105
106const OUT_TX: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
107 name: "tx",
108 ty: BuiltinParamType::Any,
109 arity: BuiltinParamArity::Required,
110 default: None,
111 description: "DataTransaction handle object.",
112}];
113
114const OUT_BOOL: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
115 name: "ok",
116 ty: BuiltinParamType::LogicalArray,
117 arity: BuiltinParamArity::Required,
118 default: None,
119 description: "Logical success/result flag.",
120}];
121
122const OUT_STRING: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
123 name: "s",
124 ty: BuiltinParamType::StringScalar,
125 arity: BuiltinParamArity::Required,
126 default: None,
127 description: "String scalar result.",
128}];
129
130const OUT_STRUCT: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
131 name: "S",
132 ty: BuiltinParamType::Any,
133 arity: BuiltinParamArity::Required,
134 default: None,
135 description: "Struct result.",
136}];
137
138const OUT_CELL: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
139 name: "C",
140 ty: BuiltinParamType::Any,
141 arity: BuiltinParamArity::Required,
142 default: None,
143 description: "Cell-array result.",
144}];
145
146const OUT_VALUE: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
147 name: "value",
148 ty: BuiltinParamType::Any,
149 arity: BuiltinParamArity::Required,
150 default: None,
151 description: "Value result.",
152}];
153
154const OUT_TENSOR: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
155 name: "X",
156 ty: BuiltinParamType::NumericArray,
157 arity: BuiltinParamArity::Required,
158 default: None,
159 description: "Numeric tensor result.",
160}];
161
162const IN_PATH_SCHEMA: [BuiltinParamDescriptor; 2] = [
163 BuiltinParamDescriptor {
164 name: "path",
165 ty: BuiltinParamType::StringScalar,
166 arity: BuiltinParamArity::Required,
167 default: None,
168 description: "Dataset path (.data).",
169 },
170 BuiltinParamDescriptor {
171 name: "schema",
172 ty: BuiltinParamType::Any,
173 arity: BuiltinParamArity::Required,
174 default: None,
175 description: "Dataset schema struct.",
176 },
177];
178
179const IN_PATH: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
180 name: "path",
181 ty: BuiltinParamType::StringScalar,
182 arity: BuiltinParamArity::Required,
183 default: None,
184 description: "Dataset path (.data).",
185}];
186
187const IN_FROM_TO: [BuiltinParamDescriptor; 2] = [
188 BuiltinParamDescriptor {
189 name: "fromPath",
190 ty: BuiltinParamType::StringScalar,
191 arity: BuiltinParamArity::Required,
192 default: None,
193 description: "Source dataset path.",
194 },
195 BuiltinParamDescriptor {
196 name: "toPath",
197 ty: BuiltinParamType::StringScalar,
198 arity: BuiltinParamArity::Required,
199 default: None,
200 description: "Destination dataset path.",
201 },
202];
203
204const IN_PATH_FORMAT_SOURCE: [BuiltinParamDescriptor; 3] = [
205 BuiltinParamDescriptor {
206 name: "path",
207 ty: BuiltinParamType::StringScalar,
208 arity: BuiltinParamArity::Required,
209 default: None,
210 description: "Dataset path (.data).",
211 },
212 BuiltinParamDescriptor {
213 name: "format",
214 ty: BuiltinParamType::StringScalar,
215 arity: BuiltinParamArity::Required,
216 default: None,
217 description: "Format token (currently 'data').",
218 },
219 BuiltinParamDescriptor {
220 name: "sourcePath",
221 ty: BuiltinParamType::StringScalar,
222 arity: BuiltinParamArity::Required,
223 default: None,
224 description: "Source dataset path.",
225 },
226];
227
228const IN_PATH_FORMAT_TARGET: [BuiltinParamDescriptor; 3] = [
229 BuiltinParamDescriptor {
230 name: "path",
231 ty: BuiltinParamType::StringScalar,
232 arity: BuiltinParamArity::Required,
233 default: None,
234 description: "Dataset path (.data).",
235 },
236 BuiltinParamDescriptor {
237 name: "format",
238 ty: BuiltinParamType::StringScalar,
239 arity: BuiltinParamArity::Required,
240 default: None,
241 description: "Format token (currently 'data').",
242 },
243 BuiltinParamDescriptor {
244 name: "targetPath",
245 ty: BuiltinParamType::StringScalar,
246 arity: BuiltinParamArity::Required,
247 default: None,
248 description: "Target dataset path.",
249 },
250];
251
252const IN_PREFIX: [BuiltinParamDescriptor; 1] = [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
260const IN_BASE: [BuiltinParamDescriptor; 1] = [BuiltinParamDescriptor {
261 name: "obj",
262 ty: BuiltinParamType::Any,
263 arity: BuiltinParamArity::Required,
264 default: None,
265 description: "Receiver object.",
266}];
267
268const IN_BASE_NAME: [BuiltinParamDescriptor; 2] = [
269 BuiltinParamDescriptor {
270 name: "obj",
271 ty: BuiltinParamType::Any,
272 arity: BuiltinParamArity::Required,
273 default: None,
274 description: "Receiver object.",
275 },
276 BuiltinParamDescriptor {
277 name: "name",
278 ty: BuiltinParamType::StringScalar,
279 arity: BuiltinParamArity::Required,
280 default: None,
281 description: "Name key/array identifier.",
282 },
283];
284
285const IN_BASE_KEY_DEFAULT: [BuiltinParamDescriptor; 3] = [
286 BuiltinParamDescriptor {
287 name: "obj",
288 ty: BuiltinParamType::Any,
289 arity: BuiltinParamArity::Required,
290 default: None,
291 description: "Receiver object.",
292 },
293 BuiltinParamDescriptor {
294 name: "key",
295 ty: BuiltinParamType::StringScalar,
296 arity: BuiltinParamArity::Required,
297 default: None,
298 description: "Attribute key.",
299 },
300 BuiltinParamDescriptor {
301 name: "defaultValue",
302 ty: BuiltinParamType::Any,
303 arity: BuiltinParamArity::Optional,
304 default: Some("0"),
305 description: "Fallback value when key is absent.",
306 },
307];
308
309const IN_BASE_KEY_VALUE: [BuiltinParamDescriptor; 3] = [
310 BuiltinParamDescriptor {
311 name: "obj",
312 ty: BuiltinParamType::Any,
313 arity: BuiltinParamArity::Required,
314 default: None,
315 description: "Receiver object.",
316 },
317 BuiltinParamDescriptor {
318 name: "key",
319 ty: BuiltinParamType::StringScalar,
320 arity: BuiltinParamArity::Required,
321 default: None,
322 description: "Attribute key.",
323 },
324 BuiltinParamDescriptor {
325 name: "value",
326 ty: BuiltinParamType::Any,
327 arity: BuiltinParamArity::Required,
328 default: None,
329 description: "Attribute value.",
330 },
331];
332
333const IN_BASE_ATTRS: [BuiltinParamDescriptor; 2] = [
334 BuiltinParamDescriptor {
335 name: "obj",
336 ty: BuiltinParamType::Any,
337 arity: BuiltinParamArity::Required,
338 default: None,
339 description: "Receiver object.",
340 },
341 BuiltinParamDescriptor {
342 name: "attrs",
343 ty: BuiltinParamType::Any,
344 arity: BuiltinParamArity::Required,
345 default: None,
346 description: "Struct of attribute updates.",
347 },
348];
349
350const IN_BASE_LABEL_REST: [BuiltinParamDescriptor; 3] = [
351 BuiltinParamDescriptor {
352 name: "obj",
353 ty: BuiltinParamType::Any,
354 arity: BuiltinParamArity::Required,
355 default: None,
356 description: "Receiver object.",
357 },
358 BuiltinParamDescriptor {
359 name: "label",
360 ty: BuiltinParamType::StringScalar,
361 arity: BuiltinParamArity::Required,
362 default: None,
363 description: "Snapshot label.",
364 },
365 BuiltinParamDescriptor {
366 name: "options",
367 ty: BuiltinParamType::Any,
368 arity: BuiltinParamArity::Variadic,
369 default: None,
370 description: "Reserved name/value options.",
371 },
372];
373
374const IN_BASE_REST: [BuiltinParamDescriptor; 2] = [
375 BuiltinParamDescriptor {
376 name: "obj",
377 ty: BuiltinParamType::Any,
378 arity: BuiltinParamArity::Required,
379 default: None,
380 description: "Receiver object.",
381 },
382 BuiltinParamDescriptor {
383 name: "options",
384 ty: BuiltinParamType::Any,
385 arity: BuiltinParamArity::Variadic,
386 default: None,
387 description: "Reserved name/value options.",
388 },
389];
390
391const IN_BASE_SLICE_OPTIONAL: [BuiltinParamDescriptor; 2] = [
392 BuiltinParamDescriptor {
393 name: "obj",
394 ty: BuiltinParamType::Any,
395 arity: BuiltinParamArity::Required,
396 default: None,
397 description: "DataArray receiver object.",
398 },
399 BuiltinParamDescriptor {
400 name: "sliceSpec",
401 ty: BuiltinParamType::Any,
402 arity: BuiltinParamArity::Optional,
403 default: None,
404 description: "Optional slice specification.",
405 },
406];
407
408const IN_BASE_VALUES: [BuiltinParamDescriptor; 2] = [
409 BuiltinParamDescriptor {
410 name: "obj",
411 ty: BuiltinParamType::Any,
412 arity: BuiltinParamArity::Required,
413 default: None,
414 description: "DataArray receiver object.",
415 },
416 BuiltinParamDescriptor {
417 name: "values",
418 ty: BuiltinParamType::Any,
419 arity: BuiltinParamArity::Required,
420 default: None,
421 description: "Full-array values payload.",
422 },
423];
424
425const IN_BASE_SLICE_VALUES: [BuiltinParamDescriptor; 3] = [
426 BuiltinParamDescriptor {
427 name: "obj",
428 ty: BuiltinParamType::Any,
429 arity: BuiltinParamArity::Required,
430 default: None,
431 description: "DataArray receiver object.",
432 },
433 BuiltinParamDescriptor {
434 name: "sliceSpec",
435 ty: BuiltinParamType::Any,
436 arity: BuiltinParamArity::Required,
437 default: None,
438 description: "Slice specification.",
439 },
440 BuiltinParamDescriptor {
441 name: "values",
442 ty: BuiltinParamType::Any,
443 arity: BuiltinParamArity::Required,
444 default: None,
445 description: "Slice values payload.",
446 },
447];
448
449const IN_BASE_NEW_SHAPE_REST: [BuiltinParamDescriptor; 3] = [
450 BuiltinParamDescriptor {
451 name: "obj",
452 ty: BuiltinParamType::Any,
453 arity: BuiltinParamArity::Required,
454 default: None,
455 description: "DataArray receiver object.",
456 },
457 BuiltinParamDescriptor {
458 name: "newShape",
459 ty: BuiltinParamType::Any,
460 arity: BuiltinParamArity::Required,
461 default: None,
462 description: "New array shape.",
463 },
464 BuiltinParamDescriptor {
465 name: "options",
466 ty: BuiltinParamType::Any,
467 arity: BuiltinParamArity::Variadic,
468 default: None,
469 description: "Reserved name/value options.",
470 },
471];
472
473const IN_BASE_VALUE_REST: [BuiltinParamDescriptor; 3] = [
474 BuiltinParamDescriptor {
475 name: "obj",
476 ty: BuiltinParamType::Any,
477 arity: BuiltinParamArity::Required,
478 default: None,
479 description: "DataArray receiver object.",
480 },
481 BuiltinParamDescriptor {
482 name: "value",
483 ty: BuiltinParamType::Any,
484 arity: BuiltinParamArity::Required,
485 default: None,
486 description: "Fill value.",
487 },
488 BuiltinParamDescriptor {
489 name: "options",
490 ty: BuiltinParamType::Any,
491 arity: BuiltinParamArity::Variadic,
492 default: None,
493 description: "Reserved name/value options.",
494 },
495];
496
497const IN_TX_WRITE: [BuiltinParamDescriptor; 5] = [
498 BuiltinParamDescriptor {
499 name: "tx",
500 ty: BuiltinParamType::Any,
501 arity: BuiltinParamArity::Required,
502 default: None,
503 description: "DataTransaction receiver.",
504 },
505 BuiltinParamDescriptor {
506 name: "arrayName",
507 ty: BuiltinParamType::StringScalar,
508 arity: BuiltinParamArity::Required,
509 default: None,
510 description: "Target array name.",
511 },
512 BuiltinParamDescriptor {
513 name: "sliceSpec",
514 ty: BuiltinParamType::Any,
515 arity: BuiltinParamArity::Required,
516 default: None,
517 description: "Slice specification.",
518 },
519 BuiltinParamDescriptor {
520 name: "values",
521 ty: BuiltinParamType::Any,
522 arity: BuiltinParamArity::Required,
523 default: None,
524 description: "Values payload.",
525 },
526 BuiltinParamDescriptor {
527 name: "options",
528 ty: BuiltinParamType::Any,
529 arity: BuiltinParamArity::Variadic,
530 default: None,
531 description: "Reserved name/value options.",
532 },
533];
534
535const IN_TX_ARRAY_SHAPE_REST: [BuiltinParamDescriptor; 4] = [
536 BuiltinParamDescriptor {
537 name: "tx",
538 ty: BuiltinParamType::Any,
539 arity: BuiltinParamArity::Required,
540 default: None,
541 description: "DataTransaction receiver.",
542 },
543 BuiltinParamDescriptor {
544 name: "arrayName",
545 ty: BuiltinParamType::StringScalar,
546 arity: BuiltinParamArity::Required,
547 default: None,
548 description: "Target array name.",
549 },
550 BuiltinParamDescriptor {
551 name: "newShape",
552 ty: BuiltinParamType::Any,
553 arity: BuiltinParamArity::Required,
554 default: None,
555 description: "New shape vector.",
556 },
557 BuiltinParamDescriptor {
558 name: "options",
559 ty: BuiltinParamType::Any,
560 arity: BuiltinParamArity::Variadic,
561 default: None,
562 description: "Reserved name/value options.",
563 },
564];
565
566const IN_TX_ARRAY_VALUE_SLICE_OPT: [BuiltinParamDescriptor; 4] = [
567 BuiltinParamDescriptor {
568 name: "tx",
569 ty: BuiltinParamType::Any,
570 arity: BuiltinParamArity::Required,
571 default: None,
572 description: "DataTransaction receiver.",
573 },
574 BuiltinParamDescriptor {
575 name: "arrayName",
576 ty: BuiltinParamType::StringScalar,
577 arity: BuiltinParamArity::Required,
578 default: None,
579 description: "Target array name.",
580 },
581 BuiltinParamDescriptor {
582 name: "value",
583 ty: BuiltinParamType::Any,
584 arity: BuiltinParamArity::Required,
585 default: None,
586 description: "Fill value.",
587 },
588 BuiltinParamDescriptor {
589 name: "sliceSpec",
590 ty: BuiltinParamType::Any,
591 arity: BuiltinParamArity::Optional,
592 default: None,
593 description: "Optional slice specification.",
594 },
595];
596
597const IN_TX_ARRAY_NAME: [BuiltinParamDescriptor; 2] = [
598 BuiltinParamDescriptor {
599 name: "tx",
600 ty: BuiltinParamType::Any,
601 arity: BuiltinParamArity::Required,
602 default: None,
603 description: "DataTransaction receiver.",
604 },
605 BuiltinParamDescriptor {
606 name: "arrayName",
607 ty: BuiltinParamType::StringScalar,
608 arity: BuiltinParamArity::Required,
609 default: None,
610 description: "Target array name.",
611 },
612];
613
614const IN_TX_ARRAY_META: [BuiltinParamDescriptor; 3] = [
615 BuiltinParamDescriptor {
616 name: "tx",
617 ty: BuiltinParamType::Any,
618 arity: BuiltinParamArity::Required,
619 default: None,
620 description: "DataTransaction receiver.",
621 },
622 BuiltinParamDescriptor {
623 name: "arrayName",
624 ty: BuiltinParamType::StringScalar,
625 arity: BuiltinParamArity::Required,
626 default: None,
627 description: "New array name.",
628 },
629 BuiltinParamDescriptor {
630 name: "meta",
631 ty: BuiltinParamType::Any,
632 arity: BuiltinParamArity::Required,
633 default: None,
634 description: "Array metadata struct.",
635 },
636];
637
638const IN_TX_COMMIT_OPTIONS: [BuiltinParamDescriptor; 2] = [
639 BuiltinParamDescriptor {
640 name: "tx",
641 ty: BuiltinParamType::Any,
642 arity: BuiltinParamArity::Required,
643 default: None,
644 description: "DataTransaction receiver.",
645 },
646 BuiltinParamDescriptor {
647 name: "options",
648 ty: BuiltinParamType::Any,
649 arity: BuiltinParamArity::Optional,
650 default: None,
651 description: "Optional struct containing an if_manifest version token.",
652 },
653];
654
655macro_rules! one_sig_descriptor {
656 ($desc:ident, $sigs:ident, $label:expr, $inputs:expr, $outputs:expr) => {
657 const $sigs: [BuiltinSignatureDescriptor; 1] = [BuiltinSignatureDescriptor {
658 label: $label,
659 inputs: $inputs,
660 outputs: $outputs,
661 }];
662 pub const $desc: BuiltinDescriptor = BuiltinDescriptor {
663 signatures: &$sigs,
664 output_mode: BuiltinOutputMode::Fixed,
665 completion_policy: BuiltinCompletionPolicy::Public,
666 errors: &DATA_DESCRIPTOR_ERRORS,
667 };
668 };
669}
670
671one_sig_descriptor!(
672 DATA_CREATE_DESCRIPTOR,
673 DATA_CREATE_SIGS,
674 "ds = data.create(path, schema)",
675 &IN_PATH_SCHEMA,
676 &OUT_DATASET
677);
678one_sig_descriptor!(
679 DATA_OPEN_DESCRIPTOR,
680 DATA_OPEN_SIGS,
681 "ds = data.open(path)",
682 &IN_PATH,
683 &OUT_DATASET
684);
685one_sig_descriptor!(
686 DATA_EXISTS_DESCRIPTOR,
687 DATA_EXISTS_SIGS,
688 "tf = data.exists(path)",
689 &IN_PATH,
690 &OUT_BOOL
691);
692one_sig_descriptor!(
693 DATA_DELETE_DESCRIPTOR,
694 DATA_DELETE_SIGS,
695 "tf = data.delete(path)",
696 &IN_PATH,
697 &OUT_BOOL
698);
699one_sig_descriptor!(
700 DATA_COPY_DESCRIPTOR,
701 DATA_COPY_SIGS,
702 "tf = data.copy(fromPath, toPath)",
703 &IN_FROM_TO,
704 &OUT_BOOL
705);
706one_sig_descriptor!(
707 DATA_MOVE_DESCRIPTOR,
708 DATA_MOVE_SIGS,
709 "tf = data.move(fromPath, toPath)",
710 &IN_FROM_TO,
711 &OUT_BOOL
712);
713one_sig_descriptor!(
714 DATA_IMPORT_DESCRIPTOR,
715 DATA_IMPORT_SIGS,
716 "ds = data.import(path, format, sourcePath)",
717 &IN_PATH_FORMAT_SOURCE,
718 &OUT_DATASET
719);
720one_sig_descriptor!(
721 DATA_EXPORT_DESCRIPTOR,
722 DATA_EXPORT_SIGS,
723 "tf = data.export(path, format, targetPath)",
724 &IN_PATH_FORMAT_TARGET,
725 &OUT_BOOL
726);
727one_sig_descriptor!(
728 DATA_LIST_DESCRIPTOR,
729 DATA_LIST_SIGS,
730 "C = data.list(prefix)",
731 &IN_PREFIX,
732 &OUT_CELL
733);
734one_sig_descriptor!(
735 DATA_INSPECT_DESCRIPTOR,
736 DATA_INSPECT_SIGS,
737 "S = data.inspect(path)",
738 &IN_PATH,
739 &OUT_STRUCT
740);
741
742const DATA_CREATE_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
743 BuiltinIntegerInputCapability {
744 name: "schema.arrays.*.shape",
745 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
746 availability: BuiltinIntegerInputAvailability::Documented,
747 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
748 notes: "Array shapes accept scalar or vector dimensions in every integer class, or exactly integral floating values, within nonnegative platform and checked total-element bounds.",
749 },
750 BuiltinIntegerInputCapability {
751 name: "schema.arrays.*.chunk",
752 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
753 availability: BuiltinIntegerInputAvailability::Documented,
754 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
755 notes: "Optional chunk shapes accept scalar or vector dimensions in every integer class, or exactly integral floating values; dimensions must be positive, platform-bounded, and rank-matched to the array shape.",
756 },
757];
758
759pub const DATA_CREATE_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
760 [BuiltinIntegerCapabilityDescriptor {
761 form: "ds = data.create(path, struct_with_integer_shape_and_chunk)",
762 inputs: &DATA_CREATE_INTEGER_INPUTS,
763 computation_domain: BuiltinIntegerComputationDomain::Structural,
764 output_class: BuiltinIntegerOutputClassRule::FunctionSpecific,
765 overflow: BuiltinIntegerOverflowRule::Error,
766 backend: BuiltinIntegerBackendRule::HostOnly,
767 overload: BuiltinIntegerOverloadKind::StructuralParameter,
768 notes: "Schema dimensions are parsed directly from authoritative typed storage; the host/filesystem constructor allocates zero-filled storage in each array's declared dtype and returns a Dataset object.",
769 }];
770
771const DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT: BuiltinIntegerAuditDescriptor =
772 BuiltinIntegerAuditDescriptor {
773 kind: BuiltinIntegerAuditKind::NotApplicable,
774 canonical_builtin: None,
775 notes: "Dataset copy, delete, existence, export, import, inspection, listing, move, and open operations accept only textual paths/formats and return objects, logicals, strings, or host metadata; opaque movement of stored integer payload bytes is owned by the DataArray persistence contract rather than an integer call form.",
776 };
777
778one_sig_descriptor!(
779 DATASET_PATH_DESCRIPTOR,
780 DATASET_PATH_SIGS,
781 "path = Dataset.path(ds)",
782 &IN_BASE,
783 &OUT_STRING
784);
785one_sig_descriptor!(
786 DATASET_ID_DESCRIPTOR,
787 DATASET_ID_SIGS,
788 "id = Dataset.id(ds)",
789 &IN_BASE,
790 &OUT_STRING
791);
792one_sig_descriptor!(
793 DATASET_VERSION_DESCRIPTOR,
794 DATASET_VERSION_SIGS,
795 "version = Dataset.version(ds)",
796 &IN_BASE,
797 &OUT_STRING
798);
799one_sig_descriptor!(
800 DATASET_ARRAYS_DESCRIPTOR,
801 DATASET_ARRAYS_SIGS,
802 "C = Dataset.arrays(ds)",
803 &IN_BASE,
804 &OUT_CELL
805);
806one_sig_descriptor!(
807 DATASET_HAS_ARRAY_DESCRIPTOR,
808 DATASET_HAS_ARRAY_SIGS,
809 "tf = Dataset.has_array(ds, name)",
810 &IN_BASE_NAME,
811 &OUT_BOOL
812);
813one_sig_descriptor!(
814 DATASET_ARRAY_DESCRIPTOR,
815 DATASET_ARRAY_SIGS,
816 "arr = Dataset.array(ds, name)",
817 &IN_BASE_NAME,
818 &OUT_ARRAY
819);
820one_sig_descriptor!(
821 DATASET_ATTRS_DESCRIPTOR,
822 DATASET_ATTRS_SIGS,
823 "S = Dataset.attrs(ds)",
824 &IN_BASE,
825 &OUT_STRUCT
826);
827one_sig_descriptor!(
828 DATASET_GET_ATTR_DESCRIPTOR,
829 DATASET_GET_ATTR_SIGS,
830 "value = Dataset.get_attr(ds, key, defaultValue)",
831 &IN_BASE_KEY_DEFAULT,
832 &OUT_VALUE
833);
834one_sig_descriptor!(
835 DATASET_SET_ATTR_DESCRIPTOR,
836 DATASET_SET_ATTR_SIGS,
837 "tf = Dataset.set_attr(ds, key, value)",
838 &IN_BASE_KEY_VALUE,
839 &OUT_BOOL
840);
841one_sig_descriptor!(
842 DATASET_SET_ATTRS_DESCRIPTOR,
843 DATASET_SET_ATTRS_SIGS,
844 "tf = Dataset.set_attrs(ds, attrs)",
845 &IN_BASE_ATTRS,
846 &OUT_BOOL
847);
848one_sig_descriptor!(
849 DATASET_BEGIN_DESCRIPTOR,
850 DATASET_BEGIN_SIGS,
851 "tx = Dataset.begin(ds, Name, Value, ...)",
852 &IN_BASE_REST,
853 &OUT_TX
854);
855one_sig_descriptor!(
856 DATASET_SNAPSHOT_DESCRIPTOR,
857 DATASET_SNAPSHOT_SIGS,
858 "snapshotPath = Dataset.snapshot(ds, label, Name, Value, ...)",
859 &IN_BASE_LABEL_REST,
860 &OUT_STRING
861);
862one_sig_descriptor!(
863 DATASET_REFRESH_DESCRIPTOR,
864 DATASET_REFRESH_SIGS,
865 "ds = Dataset.refresh(ds)",
866 &IN_BASE,
867 &OUT_DATASET
868);
869
870const DATASET_NONATTRIBUTE_INTEGER_AUDIT: BuiltinIntegerAuditDescriptor =
871 BuiltinIntegerAuditDescriptor {
872 kind: BuiltinIntegerAuditKind::NotApplicable,
873 canonical_builtin: None,
874 notes: "Dataset array lookup/listing, transaction construction, identity/path/version access, membership tests, refresh, and snapshot operations accept a Dataset receiver plus strings or reserved options; they do not accept integer values or controls, return integer-class values, perform integer arithmetic, or expose an integer backend surface.",
875 };
876
877const DATASET_ATTRS_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
878 [BuiltinIntegerInputCapability {
879 name: "stored_attribute_values",
880 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
881 availability: BuiltinIntegerInputAvailability::Documented,
882 scalar_double: BuiltinIntegerScalarDoubleRule::NotApplicable,
883 notes: "Persisted scalar attributes may originate from every integer class and retain their exact mathematical value in the JSON manifest.",
884 }];
885pub const DATASET_ATTRS_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
886 [BuiltinIntegerCapabilityDescriptor {
887 form: "S = Dataset.attrs(ds_with_integer_attributes)",
888 inputs: &DATASET_ATTRS_INTEGER_INPUTS,
889 computation_domain: BuiltinIntegerComputationDomain::Structural,
890 output_class: BuiltinIntegerOutputClassRule::FunctionSpecific,
891 overflow: BuiltinIntegerOverflowRule::NotApplicable,
892 backend: BuiltinIntegerBackendRule::HostOnly,
893 overload: BuiltinIntegerOverloadKind::Multiple,
894 notes: "Each JSON integer is decoded exactly into a canonical scalar: int64 when its value is representable there, otherwise uint64; the original narrow storage class is not encoded by JSON.",
895 }];
896
897const DATASET_GET_ATTR_STORED_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
898 [BuiltinIntegerInputCapability {
899 name: "stored_attribute",
900 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
901 availability: BuiltinIntegerInputAvailability::Documented,
902 scalar_double: BuiltinIntegerScalarDoubleRule::NotApplicable,
903 notes: "A persisted scalar attribute may originate from every integer class and retains its exact mathematical value in the JSON manifest.",
904 }];
905const DATASET_GET_ATTR_DEFAULT_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
906 [BuiltinIntegerInputCapability {
907 name: "defaultValue",
908 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
909 availability: BuiltinIntegerInputAvailability::Documented,
910 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
911 notes: "An explicit scalar default accepts every integer class and is returned unchanged when the requested key is absent.",
912 }];
913pub const DATASET_GET_ATTR_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 2] = [
914 BuiltinIntegerCapabilityDescriptor {
915 form: "value = Dataset.get_attr(ds_with_integer_attribute, key)",
916 inputs: &DATASET_GET_ATTR_STORED_INTEGER_INPUTS,
917 computation_domain: BuiltinIntegerComputationDomain::Structural,
918 output_class: BuiltinIntegerOutputClassRule::FunctionSpecific,
919 overflow: BuiltinIntegerOverflowRule::NotApplicable,
920 backend: BuiltinIntegerBackendRule::HostOnly,
921 overload: BuiltinIntegerOverloadKind::ScalarOnly,
922 notes: "A stored JSON integer is decoded exactly as int64 when representable there and otherwise as uint64; JSON does not encode the original narrow integer class.",
923 },
924 BuiltinIntegerCapabilityDescriptor {
925 form: "value = Dataset.get_attr(ds, missing_key, integer_defaultValue)",
926 inputs: &DATASET_GET_ATTR_DEFAULT_INTEGER_INPUTS,
927 computation_domain: BuiltinIntegerComputationDomain::Structural,
928 output_class: BuiltinIntegerOutputClassRule::PreserveInput,
929 overflow: BuiltinIntegerOverflowRule::NotApplicable,
930 backend: BuiltinIntegerBackendRule::HostOnly,
931 overload: BuiltinIntegerOverloadKind::ScalarOnly,
932 notes: "A missing key returns the supplied default Value directly, preserving its exact integer class and value without serialization or floating conversion.",
933 },
934];
935
936const DATASET_SET_ATTR_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
937 [BuiltinIntegerInputCapability {
938 name: "value",
939 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
940 availability: BuiltinIntegerInputAvailability::Documented,
941 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
942 notes: "A scalar attribute value accepts every integer class and is serialized as an exact JSON integer, including full-width int64 and uint64 values.",
943 }];
944pub const DATASET_SET_ATTR_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
945 [BuiltinIntegerCapabilityDescriptor {
946 form: "tf = Dataset.set_attr(ds, key, integer_value)",
947 inputs: &DATASET_SET_ATTR_INTEGER_INPUTS,
948 computation_domain: BuiltinIntegerComputationDomain::Structural,
949 output_class: BuiltinIntegerOutputClassRule::Logical,
950 overflow: BuiltinIntegerOverflowRule::NotApplicable,
951 backend: BuiltinIntegerBackendRule::HostOnly,
952 overload: BuiltinIntegerOverloadKind::ScalarOnly,
953 notes: "The scalar is emitted directly as an exact manifest JSON number without an f64 intermediary; the operation returns logical success.",
954 }];
955
956const DATASET_SET_ATTRS_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
957 [BuiltinIntegerInputCapability {
958 name: "attribute_values",
959 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
960 availability: BuiltinIntegerInputAvailability::Documented,
961 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
962 notes: "Scalar values in the attribute struct accept every integer class and are serialized as exact JSON integers, including full-width int64 and uint64 values.",
963 }];
964pub const DATASET_SET_ATTRS_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
965 [BuiltinIntegerCapabilityDescriptor {
966 form: "tf = Dataset.set_attrs(ds, struct_with_integer_values)",
967 inputs: &DATASET_SET_ATTRS_INTEGER_INPUTS,
968 computation_domain: BuiltinIntegerComputationDomain::Structural,
969 output_class: BuiltinIntegerOutputClassRule::Logical,
970 overflow: BuiltinIntegerOverflowRule::NotApplicable,
971 backend: BuiltinIntegerBackendRule::HostOnly,
972 overload: BuiltinIntegerOverloadKind::ScalarOnly,
973 notes: "Each scalar integer field is emitted directly as an exact manifest JSON number without an f64 intermediary; the operation returns logical success.",
974 }];
975
976one_sig_descriptor!(
977 DATAARRAY_NAME_DESCRIPTOR,
978 DATAARRAY_NAME_SIGS,
979 "name = DataArray.name(arr)",
980 &IN_BASE,
981 &OUT_STRING
982);
983one_sig_descriptor!(
984 DATAARRAY_DTYPE_DESCRIPTOR,
985 DATAARRAY_DTYPE_SIGS,
986 "dtype = DataArray.dtype(arr)",
987 &IN_BASE,
988 &OUT_STRING
989);
990one_sig_descriptor!(
991 DATAARRAY_SHAPE_DESCRIPTOR,
992 DATAARRAY_SHAPE_SIGS,
993 "shape = DataArray.shape(arr)",
994 &IN_BASE,
995 &OUT_TENSOR
996);
997one_sig_descriptor!(
998 DATAARRAY_RANK_DESCRIPTOR,
999 DATAARRAY_RANK_SIGS,
1000 "rank = DataArray.rank(arr)",
1001 &IN_BASE,
1002 &OUT_VALUE
1003);
1004one_sig_descriptor!(
1005 DATAARRAY_CHUNK_SHAPE_DESCRIPTOR,
1006 DATAARRAY_CHUNK_SHAPE_SIGS,
1007 "chunkShape = DataArray.chunk_shape(arr)",
1008 &IN_BASE,
1009 &OUT_TENSOR
1010);
1011one_sig_descriptor!(
1012 DATAARRAY_CODEC_DESCRIPTOR,
1013 DATAARRAY_CODEC_SIGS,
1014 "codec = DataArray.codec(arr)",
1015 &IN_BASE,
1016 &OUT_STRING
1017);
1018one_sig_descriptor!(
1019 DATAARRAY_READ_DESCRIPTOR,
1020 DATAARRAY_READ_SIGS,
1021 "X = DataArray.read(arr, sliceSpec)",
1022 &IN_BASE_SLICE_OPTIONAL,
1023 &OUT_TENSOR
1024);
1025const DATAARRAY_WRITE_SIGS: [BuiltinSignatureDescriptor; 2] = [
1026 BuiltinSignatureDescriptor {
1027 label: "tf = DataArray.write(arr, values)",
1028 inputs: &IN_BASE_VALUES,
1029 outputs: &OUT_BOOL,
1030 },
1031 BuiltinSignatureDescriptor {
1032 label: "tf = DataArray.write(arr, sliceSpec, values)",
1033 inputs: &IN_BASE_SLICE_VALUES,
1034 outputs: &OUT_BOOL,
1035 },
1036];
1037pub const DATAARRAY_WRITE_DESCRIPTOR: BuiltinDescriptor = BuiltinDescriptor {
1038 signatures: &DATAARRAY_WRITE_SIGS,
1039 output_mode: BuiltinOutputMode::Fixed,
1040 completion_policy: BuiltinCompletionPolicy::Public,
1041 errors: &DATA_DESCRIPTOR_ERRORS,
1042};
1043one_sig_descriptor!(
1044 DATAARRAY_RESIZE_DESCRIPTOR,
1045 DATAARRAY_RESIZE_SIGS,
1046 "tf = DataArray.resize(arr, newShape, Name, Value, ...)",
1047 &IN_BASE_NEW_SHAPE_REST,
1048 &OUT_BOOL
1049);
1050one_sig_descriptor!(
1051 DATAARRAY_FILL_DESCRIPTOR,
1052 DATAARRAY_FILL_SIGS,
1053 "tf = DataArray.fill(arr, value, Name, Value, ...)",
1054 &IN_BASE_VALUE_REST,
1055 &OUT_BOOL
1056);
1057
1058const DATAARRAY_METADATA_INTEGER_AUDIT: BuiltinIntegerAuditDescriptor =
1059 BuiltinIntegerAuditDescriptor {
1060 kind: BuiltinIntegerAuditKind::NotApplicable,
1061 canonical_builtin: None,
1062 notes: "DataArray name, dtype, shape, rank, chunk-shape, and codec accessors accept only a DataArray receiver object; integer storage may be described by the returned metadata but is not an integer argument, control, class-preserving value output, arithmetic form, or backend surface.",
1063 };
1064
1065const DATAARRAY_READ_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
1066 BuiltinIntegerInputCapability {
1067 name: "stored_array",
1068 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1069 availability: BuiltinIntegerInputAvailability::Documented,
1070 scalar_double: BuiltinIntegerScalarDoubleRule::NotApplicable,
1071 notes: "A DataArray whose declared dtype is any built-in integer class is decoded directly into authoritative same-class real or paired-complex tensor storage.",
1072 },
1073 BuiltinIntegerInputCapability {
1074 name: "sliceSpec",
1075 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1076 availability: BuiltinIntegerInputAvailability::Documented,
1077 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1078 notes: "Scalar indices and two-element inclusive ranges accept all eight integer classes or exactly integral floating values.",
1079 },
1080];
1081pub const DATAARRAY_READ_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1082 [BuiltinIntegerCapabilityDescriptor {
1083 form: "X = DataArray.read(integer_array, integer_sliceSpec)",
1084 inputs: &DATAARRAY_READ_INTEGER_INPUTS,
1085 computation_domain: BuiltinIntegerComputationDomain::Structural,
1086 output_class: BuiltinIntegerOutputClassRule::PreserveInput,
1087 overflow: BuiltinIntegerOverflowRule::NotApplicable,
1088 backend: BuiltinIntegerBackendRule::HostOnly,
1089 overload: BuiltinIntegerOverloadKind::Multiple,
1090 notes: "Full and sliced reads preserve the declared integer dtype, real or paired-complex storage, exact values, selected shape, and column-major order; dataset persistence is host/filesystem I/O.",
1091 }];
1092
1093const DATAARRAY_WRITE_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
1094 BuiltinIntegerInputCapability {
1095 name: "values",
1096 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1097 availability: BuiltinIntegerInputAvailability::Documented,
1098 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1099 notes: "Scalar or tensor values accept every integer class in real or paired-complex storage and are stored exactly when their class matches the declared array dtype.",
1100 },
1101 BuiltinIntegerInputCapability {
1102 name: "sliceSpec",
1103 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1104 availability: BuiltinIntegerInputAvailability::Documented,
1105 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1106 notes: "Optional scalar indices and two-element inclusive ranges accept all eight integer classes or exactly integral floating values.",
1107 },
1108];
1109pub const DATAARRAY_WRITE_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1110 [BuiltinIntegerCapabilityDescriptor {
1111 form: "tf = DataArray.write(arr, integer_values) or DataArray.write(arr, integer_sliceSpec, integer_values)",
1112 inputs: &DATAARRAY_WRITE_INTEGER_INPUTS,
1113 computation_domain: BuiltinIntegerComputationDomain::Structural,
1114 output_class: BuiltinIntegerOutputClassRule::Logical,
1115 overflow: BuiltinIntegerOverflowRule::FunctionSpecific,
1116 backend: BuiltinIntegerBackendRule::HostOnly,
1117 overload: BuiltinIntegerOverloadKind::Multiple,
1118 notes: "Writes preserve exact same-class real or paired-complex integer storage; unlike source and declared dtypes use the declared numeric cast contract before chunk serialization, and shape or complexity mismatches reject.",
1119 }];
1120
1121const DATAARRAY_RESIZE_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
1122 [BuiltinIntegerInputCapability {
1123 name: "newShape",
1124 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1125 availability: BuiltinIntegerInputAvailability::Documented,
1126 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1127 notes: "Shape scalars and vectors accept all eight integer classes or exactly integral floating values within nonnegative platform and checked total-element bounds.",
1128 }];
1129pub const DATAARRAY_RESIZE_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1130 [BuiltinIntegerCapabilityDescriptor {
1131 form: "tf = DataArray.resize(arr, integer_newShape)",
1132 inputs: &DATAARRAY_RESIZE_INTEGER_INPUTS,
1133 computation_domain: BuiltinIntegerComputationDomain::Structural,
1134 output_class: BuiltinIntegerOutputClassRule::Logical,
1135 overflow: BuiltinIntegerOverflowRule::Error,
1136 backend: BuiltinIntegerBackendRule::HostOnly,
1137 overload: BuiltinIntegerOverloadKind::StructuralParameter,
1138 notes: "The exact shape is validated against platform bounds, persisted as metadata, and used to recreate zero-filled storage in the array's declared dtype.",
1139 }];
1140
1141const DATAARRAY_FILL_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
1142 [BuiltinIntegerInputCapability {
1143 name: "value",
1144 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1145 availability: BuiltinIntegerInputAvailability::Documented,
1146 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1147 notes: "The scalar fill value accepts every integer class and is exact when its class matches the declared array dtype.",
1148 }];
1149pub const DATAARRAY_FILL_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1150 [BuiltinIntegerCapabilityDescriptor {
1151 form: "tf = DataArray.fill(arr, integer_value)",
1152 inputs: &DATAARRAY_FILL_INTEGER_INPUTS,
1153 computation_domain: BuiltinIntegerComputationDomain::Structural,
1154 output_class: BuiltinIntegerOutputClassRule::Logical,
1155 overflow: BuiltinIntegerOverflowRule::FunctionSpecific,
1156 backend: BuiltinIntegerBackendRule::HostOnly,
1157 overload: BuiltinIntegerOverloadKind::ScalarOnly,
1158 notes: "A scalar is cast once to the declared dtype and that exact typed value is repeated; nonscalar fills reject.",
1159 }];
1160
1161const DATATX_LIFECYCLE_INTEGER_AUDIT: BuiltinIntegerAuditDescriptor =
1162 BuiltinIntegerAuditDescriptor {
1163 kind: BuiltinIntegerAuditKind::NotApplicable,
1164 canonical_builtin: None,
1165 notes: "DataTransaction id, status, commit, abort, and delete-array operations accept a transaction receiver plus, where applicable, string metadata or an array name; they do not accept integer values or controls, return integer-class values, perform integer arithmetic, or expose an integer backend surface.",
1166 };
1167
1168const DATATX_WRITE_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
1169 BuiltinIntegerInputCapability {
1170 name: "values",
1171 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1172 availability: BuiltinIntegerInputAvailability::Documented,
1173 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1174 notes: "Queued scalar or tensor values accept every integer class and are stored exactly at commit when their class matches the declared array dtype.",
1175 },
1176 BuiltinIntegerInputCapability {
1177 name: "sliceSpec",
1178 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1179 availability: BuiltinIntegerInputAvailability::Documented,
1180 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1181 notes: "Queued scalar indices and two-element inclusive ranges accept all eight integer classes or exactly integral floating values.",
1182 },
1183];
1184pub const DATATX_WRITE_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1185 [BuiltinIntegerCapabilityDescriptor {
1186 form: "tf = DataTransaction.write(tx, arrayName, integer_sliceSpec, integer_values)",
1187 inputs: &DATATX_WRITE_INTEGER_INPUTS,
1188 computation_domain: BuiltinIntegerComputationDomain::Structural,
1189 output_class: BuiltinIntegerOutputClassRule::Logical,
1190 overflow: BuiltinIntegerOverflowRule::FunctionSpecific,
1191 backend: BuiltinIntegerBackendRule::HostOnly,
1192 overload: BuiltinIntegerOverloadKind::Multiple,
1193 notes: "The transaction retains native integer values and exact slice controls until commit; commit applies the same declared-dtype cast, shape validation, and chunk serialization contract as DataArray.write.",
1194 }];
1195
1196const DATATX_SET_ATTR_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
1197 [BuiltinIntegerInputCapability {
1198 name: "value",
1199 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1200 availability: BuiltinIntegerInputAvailability::Documented,
1201 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1202 notes: "A scalar attribute value accepts every integer class and is serialized as an exact JSON integer, including full-width int64 and uint64 values.",
1203 }];
1204pub const DATATX_SET_ATTR_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1205 [BuiltinIntegerCapabilityDescriptor {
1206 form: "tf = DataTransaction.set_attr(tx, key, integer_value)",
1207 inputs: &DATATX_SET_ATTR_INTEGER_INPUTS,
1208 computation_domain: BuiltinIntegerComputationDomain::Structural,
1209 output_class: BuiltinIntegerOutputClassRule::Logical,
1210 overflow: BuiltinIntegerOverflowRule::NotApplicable,
1211 backend: BuiltinIntegerBackendRule::HostOnly,
1212 overload: BuiltinIntegerOverloadKind::ScalarOnly,
1213 notes: "The queued scalar retains its native integer class through the transaction and is emitted as an exact manifest JSON number at commit without an f64 intermediary.",
1214 }];
1215
1216const DATATX_SET_ATTRS_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
1217 [BuiltinIntegerInputCapability {
1218 name: "attribute_values",
1219 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1220 availability: BuiltinIntegerInputAvailability::Documented,
1221 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1222 notes: "Scalar values in the attribute struct accept every integer class and are serialized as exact JSON integers, including full-width int64 and uint64 values.",
1223 }];
1224pub const DATATX_SET_ATTRS_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1225 [BuiltinIntegerCapabilityDescriptor {
1226 form: "tf = DataTransaction.set_attrs(tx, struct_with_integer_values)",
1227 inputs: &DATATX_SET_ATTRS_INTEGER_INPUTS,
1228 computation_domain: BuiltinIntegerComputationDomain::Structural,
1229 output_class: BuiltinIntegerOutputClassRule::Logical,
1230 overflow: BuiltinIntegerOverflowRule::NotApplicable,
1231 backend: BuiltinIntegerBackendRule::HostOnly,
1232 overload: BuiltinIntegerOverloadKind::ScalarOnly,
1233 notes: "Each queued scalar integer field retains its native class and is emitted as an exact manifest JSON number at commit without an f64 intermediary.",
1234 }];
1235
1236const DATATX_RESIZE_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 1] =
1237 [BuiltinIntegerInputCapability {
1238 name: "newShape",
1239 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1240 availability: BuiltinIntegerInputAvailability::Documented,
1241 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1242 notes: "Queued shape scalars and vectors accept all eight integer classes or exactly integral floating values within nonnegative platform and checked total-element bounds.",
1243 }];
1244pub const DATATX_RESIZE_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1245 [BuiltinIntegerCapabilityDescriptor {
1246 form: "tf = DataTransaction.resize(tx, arrayName, integer_newShape)",
1247 inputs: &DATATX_RESIZE_INTEGER_INPUTS,
1248 computation_domain: BuiltinIntegerComputationDomain::Structural,
1249 output_class: BuiltinIntegerOutputClassRule::Logical,
1250 overflow: BuiltinIntegerOverflowRule::Error,
1251 backend: BuiltinIntegerBackendRule::HostOnly,
1252 overload: BuiltinIntegerOverloadKind::StructuralParameter,
1253 notes: "The exact shape is validated when queued, retained as platform-sized metadata, and applied at commit by recreating zero-filled storage in the array's declared dtype.",
1254 }];
1255
1256const DATATX_FILL_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
1257 BuiltinIntegerInputCapability {
1258 name: "value",
1259 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1260 availability: BuiltinIntegerInputAvailability::Documented,
1261 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1262 notes: "The queued scalar fill value accepts every integer class and is exact at commit when its class matches the declared array dtype.",
1263 },
1264 BuiltinIntegerInputCapability {
1265 name: "sliceSpec",
1266 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1267 availability: BuiltinIntegerInputAvailability::Documented,
1268 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1269 notes: "Optional queued scalar indices and two-element inclusive ranges accept all eight integer classes or exactly integral floating values.",
1270 },
1271];
1272pub const DATATX_FILL_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1273 [BuiltinIntegerCapabilityDescriptor {
1274 form: "tf = DataTransaction.fill(tx, arrayName, integer_value, integer_sliceSpec)",
1275 inputs: &DATATX_FILL_INTEGER_INPUTS,
1276 computation_domain: BuiltinIntegerComputationDomain::Structural,
1277 output_class: BuiltinIntegerOutputClassRule::Logical,
1278 overflow: BuiltinIntegerOverflowRule::FunctionSpecific,
1279 backend: BuiltinIntegerBackendRule::HostOnly,
1280 overload: BuiltinIntegerOverloadKind::Multiple,
1281 notes: "The transaction retains the exact scalar and optional slice until commit; commit casts once to the declared dtype and repeats that typed value.",
1282 }];
1283
1284const DATATX_CREATE_ARRAY_INTEGER_INPUTS: [BuiltinIntegerInputCapability; 2] = [
1285 BuiltinIntegerInputCapability {
1286 name: "meta.shape",
1287 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1288 availability: BuiltinIntegerInputAvailability::Documented,
1289 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1290 notes: "The schema shape accepts all eight integer classes or exactly integral floating values within nonnegative platform and checked total-element bounds.",
1291 },
1292 BuiltinIntegerInputCapability {
1293 name: "meta.chunk",
1294 classes: &crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES,
1295 availability: BuiltinIntegerInputAvailability::Documented,
1296 scalar_double: BuiltinIntegerScalarDoubleRule::Allowed,
1297 notes: "The optional chunk shape accepts all eight integer classes or exactly integral floating values; dimensions must be positive, platform-bounded, and rank-matched to meta.shape.",
1298 },
1299];
1300pub const DATATX_CREATE_ARRAY_INTEGER_CAPABILITIES: [BuiltinIntegerCapabilityDescriptor; 1] =
1301 [BuiltinIntegerCapabilityDescriptor {
1302 form: "tf = DataTransaction.create_array(tx, arrayName, struct('shape', integer_shape, 'chunk', integer_chunk))",
1303 inputs: &DATATX_CREATE_ARRAY_INTEGER_INPUTS,
1304 computation_domain: BuiltinIntegerComputationDomain::Structural,
1305 output_class: BuiltinIntegerOutputClassRule::Logical,
1306 overflow: BuiltinIntegerOverflowRule::Error,
1307 backend: BuiltinIntegerBackendRule::HostOnly,
1308 overload: BuiltinIntegerOverloadKind::StructuralParameter,
1309 notes: "Exact shape and chunk dimensions are validated when queued, retained as platform-sized metadata, and used at commit to create storage in the schema's declared dtype.",
1310 }];
1311
1312one_sig_descriptor!(
1313 DATATX_ID_DESCRIPTOR,
1314 DATATX_ID_SIGS,
1315 "id = DataTransaction.id(tx)",
1316 &IN_BASE,
1317 &OUT_STRING
1318);
1319one_sig_descriptor!(
1320 DATATX_WRITE_DESCRIPTOR,
1321 DATATX_WRITE_SIGS_1,
1322 "tf = DataTransaction.write(tx, arrayName, sliceSpec, values, Name, Value, ...)",
1323 &IN_TX_WRITE,
1324 &OUT_BOOL
1325);
1326one_sig_descriptor!(
1327 DATATX_SET_ATTR_DESCRIPTOR,
1328 DATATX_SET_ATTR_SIGS,
1329 "tf = DataTransaction.set_attr(tx, key, value)",
1330 &IN_BASE_KEY_VALUE,
1331 &OUT_BOOL
1332);
1333one_sig_descriptor!(
1334 DATATX_SET_ATTRS_DESCRIPTOR,
1335 DATATX_SET_ATTRS_SIGS,
1336 "tf = DataTransaction.set_attrs(tx, attrs)",
1337 &IN_BASE_ATTRS,
1338 &OUT_BOOL
1339);
1340one_sig_descriptor!(
1341 DATATX_RESIZE_DESCRIPTOR,
1342 DATATX_RESIZE_SIGS,
1343 "tf = DataTransaction.resize(tx, arrayName, newShape, Name, Value, ...)",
1344 &IN_TX_ARRAY_SHAPE_REST,
1345 &OUT_BOOL
1346);
1347one_sig_descriptor!(
1348 DATATX_FILL_DESCRIPTOR,
1349 DATATX_FILL_SIGS,
1350 "tf = DataTransaction.fill(tx, arrayName, value, sliceSpec)",
1351 &IN_TX_ARRAY_VALUE_SLICE_OPT,
1352 &OUT_BOOL
1353);
1354one_sig_descriptor!(
1355 DATATX_DELETE_ARRAY_DESCRIPTOR,
1356 DATATX_DELETE_ARRAY_SIGS,
1357 "tf = DataTransaction.delete_array(tx, arrayName)",
1358 &IN_TX_ARRAY_NAME,
1359 &OUT_BOOL
1360);
1361one_sig_descriptor!(
1362 DATATX_CREATE_ARRAY_DESCRIPTOR,
1363 DATATX_CREATE_ARRAY_SIGS,
1364 "tf = DataTransaction.create_array(tx, arrayName, meta)",
1365 &IN_TX_ARRAY_META,
1366 &OUT_BOOL
1367);
1368one_sig_descriptor!(
1369 DATATX_COMMIT_DESCRIPTOR,
1370 DATATX_COMMIT_SIGS,
1371 "tf = DataTransaction.commit(tx, options)",
1372 &IN_TX_COMMIT_OPTIONS,
1373 &OUT_BOOL
1374);
1375one_sig_descriptor!(
1376 COMMIT_ALIAS_DESCRIPTOR,
1377 COMMIT_ALIAS_SIGS,
1378 "tf = commit(tx, options)",
1379 &IN_TX_COMMIT_OPTIONS,
1380 &OUT_BOOL
1381);
1382one_sig_descriptor!(
1383 DATATX_ABORT_DESCRIPTOR,
1384 DATATX_ABORT_SIGS,
1385 "tf = DataTransaction.abort(tx)",
1386 &IN_BASE,
1387 &OUT_BOOL
1388);
1389one_sig_descriptor!(
1390 DATATX_STATUS_DESCRIPTOR,
1391 DATATX_STATUS_SIGS,
1392 "status = DataTransaction.status(tx)",
1393 &IN_BASE,
1394 &OUT_STRING
1395);
1396
1397#[runtime_builtin(
1398 name = "data.create",
1399 category = "io/data",
1400 summary = "Create a typed dataset at a .data path.",
1401 keywords = "data,dataset,create,persistence",
1402 sink = true,
1403 type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1404 descriptor(crate::builtins::io::data::DATA_CREATE_DESCRIPTOR),
1405 integer_capabilities(crate::builtins::io::data::DATA_CREATE_INTEGER_CAPABILITIES),
1406 builtin_path = "crate::builtins::io::data"
1407)]
1408async fn data_create_builtin(path: Value, schema: Value) -> BuiltinResult<Value> {
1409 let path = parse_string(&path, "data.create path")?;
1410 let root = dataset_root(&path);
1411 let schema = parse_schema(&schema)?;
1412 let now = now_rfc3339();
1413 let mut arrays = BTreeMap::new();
1414 for (name, mut meta) in schema.arrays {
1415 let payload = DataArrayPayload::zeros(meta.dtype.clone(), meta.shape.clone())?;
1416 let (payload_path, chunk_index_path) =
1417 write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
1418 meta.data_path = make_rel_data_path(&root, &payload_path)?;
1419 meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
1420 arrays.insert(name, meta);
1421 }
1422
1423 let manifest = DataManifest {
1424 schema_version: 1,
1425 format: "runmat-data".to_string(),
1426 dataset_id: crate::data::new_dataset_id(),
1427 name: root.file_name().map(|v| v.to_string_lossy().to_string()),
1428 created_at: now.clone(),
1429 updated_at: now,
1430 arrays,
1431 attrs: BTreeMap::new(),
1432 txn_sequence: 0,
1433 };
1434 write_manifest_async(&root, &manifest).await?;
1435 Ok(dataset_object(&path, &manifest))
1436}
1437
1438#[runtime_builtin(
1439 name = "data.open",
1440 category = "io/data",
1441 summary = "Open a dataset handle from a .data path.",
1442 keywords = "data,dataset,open,persistence",
1443 type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1444 descriptor(crate::builtins::io::data::DATA_OPEN_DESCRIPTOR),
1445 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1446 builtin_path = "crate::builtins::io::data"
1447)]
1448async fn data_open_builtin(path: Value) -> BuiltinResult<Value> {
1449 let path = parse_string(&path, "data.open path")?;
1450 let root = dataset_root(&path);
1451 let manifest = read_manifest_async(&root).await?;
1452 let mut ds = dataset_object(&path, &manifest);
1453 hydrate_dataset_descriptor_async(&path, &mut ds).await;
1454 Ok(ds)
1455}
1456
1457#[runtime_builtin(
1458 name = "data.exists",
1459 category = "io/data",
1460 summary = "Check if dataset exists.",
1461 keywords = "data,dataset,exists",
1462 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1463 descriptor(crate::builtins::io::data::DATA_EXISTS_DESCRIPTOR),
1464 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1465 builtin_path = "crate::builtins::io::data"
1466)]
1467async fn data_exists_builtin(path: Value) -> BuiltinResult<Value> {
1468 let path = parse_string(&path, "data.exists path")?;
1469 let root = dataset_root(&path);
1470 let exists = runmat_filesystem::metadata_async(manifest_path(&root))
1471 .await
1472 .is_ok();
1473 Ok(Value::Bool(exists))
1474}
1475
1476#[runtime_builtin(
1477 name = "data.delete",
1478 category = "io/data",
1479 summary = "Delete a dataset path.",
1480 keywords = "data,dataset,delete",
1481 sink = true,
1482 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1483 descriptor(crate::builtins::io::data::DATA_DELETE_DESCRIPTOR),
1484 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1485 builtin_path = "crate::builtins::io::data"
1486)]
1487async fn data_delete_builtin(path: Value) -> BuiltinResult<Value> {
1488 let path = parse_string(&path, "data.delete path")?;
1489 let root = dataset_root(&path);
1490 runmat_filesystem::remove_dir_all_async(&root)
1491 .await
1492 .map_err(|err| {
1493 data_error(format!(
1494 "data.delete: failed to remove '{}': {err}",
1495 root.display()
1496 ))
1497 })?;
1498 Ok(Value::Bool(true))
1499}
1500
1501#[runtime_builtin(
1502 name = "data.copy",
1503 category = "io/data",
1504 summary = "Copy dataset to new path.",
1505 keywords = "data,dataset,copy",
1506 sink = true,
1507 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1508 descriptor(crate::builtins::io::data::DATA_COPY_DESCRIPTOR),
1509 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1510 builtin_path = "crate::builtins::io::data"
1511)]
1512async fn data_copy_builtin(from_path: Value, to_path: Value) -> BuiltinResult<Value> {
1513 let from = parse_string(&from_path, "data.copy fromPath")?;
1514 let to = parse_string(&to_path, "data.copy toPath")?;
1515 copy_dir_recursive(&dataset_root(&from), &dataset_root(&to)).await?;
1516 Ok(Value::Bool(true))
1517}
1518
1519#[runtime_builtin(
1520 name = "data.move",
1521 category = "io/data",
1522 summary = "Move dataset to new path.",
1523 keywords = "data,dataset,move",
1524 sink = true,
1525 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1526 descriptor(crate::builtins::io::data::DATA_MOVE_DESCRIPTOR),
1527 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1528 builtin_path = "crate::builtins::io::data"
1529)]
1530async fn data_move_builtin(from_path: Value, to_path: Value) -> BuiltinResult<Value> {
1531 let from = parse_string(&from_path, "data.move fromPath")?;
1532 let to = parse_string(&to_path, "data.move toPath")?;
1533 runmat_filesystem::rename_async(dataset_root(&from), dataset_root(&to))
1534 .await
1535 .map_err(|err| {
1536 data_error(format!(
1537 "data.move: failed to move dataset '{from}' -> '{to}': {err}"
1538 ))
1539 })?;
1540 Ok(Value::Bool(true))
1541}
1542
1543#[runtime_builtin(
1544 name = "data.import",
1545 category = "io/data",
1546 summary = "Import an existing dataset file path.",
1547 keywords = "data,dataset,import",
1548 sink = true,
1549 type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1550 descriptor(crate::builtins::io::data::DATA_IMPORT_DESCRIPTOR),
1551 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1552 builtin_path = "crate::builtins::io::data"
1553)]
1554async fn data_import_builtin(
1555 path: Value,
1556 format: Value,
1557 source_path: Value,
1558) -> BuiltinResult<Value> {
1559 let path = parse_string(&path, "data.import path")?;
1560 let format = parse_string(&format, "data.import format")?;
1561 if !format.eq_ignore_ascii_case("data") {
1562 return Err(data_error(
1563 "data.import currently supports only format='data'",
1564 ));
1565 }
1566 let source_path = parse_string(&source_path, "data.import sourcePath")?;
1567 copy_dir_recursive(&dataset_root(&source_path), &dataset_root(&path)).await?;
1568 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1569 Ok(dataset_object(&path, &manifest))
1570}
1571
1572#[runtime_builtin(
1573 name = "data.export",
1574 category = "io/data",
1575 summary = "Export dataset to target path.",
1576 keywords = "data,dataset,export",
1577 sink = true,
1578 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1579 descriptor(crate::builtins::io::data::DATA_EXPORT_DESCRIPTOR),
1580 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1581 builtin_path = "crate::builtins::io::data"
1582)]
1583async fn data_export_builtin(
1584 path: Value,
1585 format: Value,
1586 target_path: Value,
1587) -> BuiltinResult<Value> {
1588 let path = parse_string(&path, "data.export path")?;
1589 let format = parse_string(&format, "data.export format")?;
1590 if !format.eq_ignore_ascii_case("data") {
1591 return Err(data_error(
1592 "data.export currently supports only format='data'",
1593 ));
1594 }
1595 let target_path = parse_string(&target_path, "data.export targetPath")?;
1596 copy_dir_recursive(&dataset_root(&path), &dataset_root(&target_path)).await?;
1597 Ok(Value::Bool(true))
1598}
1599
1600#[runtime_builtin(
1601 name = "data.list",
1602 category = "io/data",
1603 summary = "List dataset paths under a prefix.",
1604 keywords = "data,dataset,list",
1605 type_resolver(crate::builtins::io::type_resolvers::data_cell_string_type),
1606 descriptor(crate::builtins::io::data::DATA_LIST_DESCRIPTOR),
1607 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1608 builtin_path = "crate::builtins::io::data"
1609)]
1610async fn data_list_builtin(path_prefix: Value) -> BuiltinResult<Value> {
1611 let prefix = parse_string(&path_prefix, "data.list prefix")?;
1612 let root = PathBuf::from(prefix);
1613 let entries = runmat_filesystem::read_dir_async(&root)
1614 .await
1615 .map_err(|err| {
1616 data_error(format!(
1617 "data.list: failed to read '{}': {err}",
1618 root.display()
1619 ))
1620 })?;
1621 let mut values = Vec::new();
1622 for entry in entries {
1623 if !entry.is_dir() {
1624 continue;
1625 }
1626 let candidate = entry.path();
1627 if candidate.extension().and_then(|s| s.to_str()) != Some("data") {
1628 continue;
1629 }
1630 if runmat_filesystem::metadata_async(candidate.join("manifest.json"))
1631 .await
1632 .is_ok()
1633 {
1634 values.push(Value::String(candidate.to_string_lossy().to_string()));
1635 }
1636 }
1637 let cols = values.len();
1638 make_cell(values, 1, cols).map_err(data_error)
1639}
1640
1641#[runtime_builtin(
1642 name = "data.inspect",
1643 category = "io/data",
1644 summary = "Inspect dataset metadata and schema fields.",
1645 keywords = "data,dataset,inspect,schema",
1646 type_resolver(crate::builtins::io::type_resolvers::data_struct_type),
1647 descriptor(crate::builtins::io::data::DATA_INSPECT_DESCRIPTOR),
1648 integer_audit(crate::builtins::io::data::DATA_NAMESPACE_LIFECYCLE_INTEGER_AUDIT),
1649 builtin_path = "crate::builtins::io::data"
1650)]
1651async fn data_inspect_builtin(path: Value) -> BuiltinResult<Value> {
1652 let path = parse_string(&path, "data.inspect path")?;
1653 let root = dataset_root(&path);
1654 let manifest = read_manifest_async(&root).await?;
1655 let mut out = StructValue::new();
1656 out.fields.insert("path".to_string(), Value::String(path));
1657 out.fields
1658 .insert("id".to_string(), Value::String(manifest.dataset_id));
1659 out.fields.insert(
1660 "arrayCount".to_string(),
1661 Value::Num(manifest.arrays.len() as f64),
1662 );
1663 out.fields
1664 .insert("updatedAt".to_string(), Value::String(manifest.updated_at));
1665 Ok(Value::Struct(out))
1666}
1667
1668#[runtime_builtin(
1669 name = "Dataset.path",
1670 category = "io/data",
1671 summary = "Return dataset path.",
1672 keywords = "dataset,path",
1673 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1674 descriptor(crate::builtins::io::data::DATASET_PATH_DESCRIPTOR),
1675 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1676 builtin_path = "crate::builtins::io::data"
1677)]
1678async fn dataset_path_builtin(base: Value) -> BuiltinResult<Value> {
1679 let obj = as_object(&base, "Dataset.path")?;
1680 Ok(get_object_prop(obj, "__data_path")?.clone())
1681}
1682
1683#[runtime_builtin(
1684 name = "Dataset.id",
1685 category = "io/data",
1686 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1687 descriptor(crate::builtins::io::data::DATASET_ID_DESCRIPTOR),
1688 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1689 builtin_path = "crate::builtins::io::data"
1690)]
1691async fn dataset_id_builtin(base: Value) -> BuiltinResult<Value> {
1692 let obj = as_object(&base, "Dataset.id")?;
1693 Ok(get_object_prop(obj, "__data_id")?.clone())
1694}
1695
1696#[runtime_builtin(
1697 name = "Dataset.version",
1698 category = "io/data",
1699 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1700 descriptor(crate::builtins::io::data::DATASET_VERSION_DESCRIPTOR),
1701 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1702 builtin_path = "crate::builtins::io::data"
1703)]
1704async fn dataset_version_builtin(base: Value) -> BuiltinResult<Value> {
1705 let obj = as_object(&base, "Dataset.version")?;
1706 Ok(get_object_prop(obj, "__data_version")?.clone())
1707}
1708
1709#[runtime_builtin(
1710 name = "Dataset.arrays",
1711 category = "io/data",
1712 type_resolver(crate::builtins::io::type_resolvers::data_cell_string_type),
1713 descriptor(crate::builtins::io::data::DATASET_ARRAYS_DESCRIPTOR),
1714 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1715 builtin_path = "crate::builtins::io::data"
1716)]
1717async fn dataset_arrays_builtin(base: Value) -> BuiltinResult<Value> {
1718 let path = dataset_path_from_object(&base, "Dataset.arrays")?;
1719 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1720 let values: Vec<Value> = manifest
1721 .arrays
1722 .keys()
1723 .map(|k| Value::String(k.clone()))
1724 .collect();
1725 make_cell(values.clone(), 1, values.len()).map_err(data_error)
1726}
1727
1728#[runtime_builtin(
1729 name = "Dataset.has_array",
1730 category = "io/data",
1731 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1732 descriptor(crate::builtins::io::data::DATASET_HAS_ARRAY_DESCRIPTOR),
1733 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1734 builtin_path = "crate::builtins::io::data"
1735)]
1736async fn dataset_has_array_builtin(base: Value, name: Value) -> BuiltinResult<Value> {
1737 let path = dataset_path_from_object(&base, "Dataset.has_array")?;
1738 let name = parse_string(&name, "Dataset.has_array name")?;
1739 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1740 Ok(Value::Bool(manifest.arrays.contains_key(&name)))
1741}
1742
1743#[runtime_builtin(
1744 name = "Dataset.array",
1745 category = "io/data",
1746 type_resolver(crate::builtins::io::type_resolvers::data_array_type),
1747 descriptor(crate::builtins::io::data::DATASET_ARRAY_DESCRIPTOR),
1748 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1749 builtin_path = "crate::builtins::io::data"
1750)]
1751async fn dataset_array_builtin(base: Value, name: Value) -> BuiltinResult<Value> {
1752 let path = dataset_path_from_object(&base, "Dataset.array")?;
1753 let name = parse_string(&name, "Dataset.array name")?;
1754 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1755 if !manifest.arrays.contains_key(&name) {
1756 return Err(data_error(format!(
1757 "Dataset.array: array '{name}' not found"
1758 )));
1759 }
1760 Ok(array_object(&path, &name))
1761}
1762
1763#[runtime_builtin(
1764 name = "Dataset.attrs",
1765 category = "io/data",
1766 type_resolver(crate::builtins::io::type_resolvers::data_struct_type),
1767 descriptor(crate::builtins::io::data::DATASET_ATTRS_DESCRIPTOR),
1768 integer_capabilities(crate::builtins::io::data::DATASET_ATTRS_INTEGER_CAPABILITIES),
1769 builtin_path = "crate::builtins::io::data"
1770)]
1771async fn dataset_attrs_builtin(base: Value) -> BuiltinResult<Value> {
1772 let path = dataset_path_from_object(&base, "Dataset.attrs")?;
1773 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1774 Ok(attrs_to_struct(&manifest.attrs))
1775}
1776
1777#[runtime_builtin(
1778 name = "Dataset.get_attr",
1779 category = "io/data",
1780 type_resolver(crate::builtins::io::type_resolvers::data_unknown_type),
1781 descriptor(crate::builtins::io::data::DATASET_GET_ATTR_DESCRIPTOR),
1782 integer_capabilities(crate::builtins::io::data::DATASET_GET_ATTR_INTEGER_CAPABILITIES),
1783 builtin_path = "crate::builtins::io::data"
1784)]
1785async fn dataset_get_attr_builtin(
1786 base: Value,
1787 key: Value,
1788 rest: Vec<Value>,
1789) -> BuiltinResult<Value> {
1790 let path = dataset_path_from_object(&base, "Dataset.get_attr")?;
1791 let key = parse_string(&key, "Dataset.get_attr key")?;
1792 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1793 if let Some(value) = manifest.attrs.get(&key) {
1794 return Ok(json_to_value(value));
1795 }
1796 Ok(rest.first().cloned().unwrap_or(Value::Num(0.0)))
1797}
1798
1799#[runtime_builtin(
1800 name = "Dataset.set_attr",
1801 category = "io/data",
1802 sink = true,
1803 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1804 descriptor(crate::builtins::io::data::DATASET_SET_ATTR_DESCRIPTOR),
1805 integer_capabilities(crate::builtins::io::data::DATASET_SET_ATTR_INTEGER_CAPABILITIES),
1806 builtin_path = "crate::builtins::io::data"
1807)]
1808async fn dataset_set_attr_builtin(base: Value, key: Value, value: Value) -> BuiltinResult<Value> {
1809 let path = dataset_path_from_object(&base, "Dataset.set_attr")?;
1810 let key = parse_string(&key, "Dataset.set_attr key")?;
1811 let root = dataset_root(&path);
1812 let mut manifest = read_manifest_async(&root).await?;
1813 manifest.attrs.insert(key, value_to_json(&value));
1814 manifest.updated_at = now_rfc3339();
1815 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1816 write_manifest_async(&root, &manifest).await?;
1817 Ok(Value::Bool(true))
1818}
1819
1820#[runtime_builtin(
1821 name = "Dataset.set_attrs",
1822 category = "io/data",
1823 sink = true,
1824 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
1825 descriptor(crate::builtins::io::data::DATASET_SET_ATTRS_DESCRIPTOR),
1826 integer_capabilities(crate::builtins::io::data::DATASET_SET_ATTRS_INTEGER_CAPABILITIES),
1827 builtin_path = "crate::builtins::io::data"
1828)]
1829async fn dataset_set_attrs_builtin(base: Value, attrs: Value) -> BuiltinResult<Value> {
1830 let path = dataset_path_from_object(&base, "Dataset.set_attrs")?;
1831 let Value::Struct(incoming) = attrs else {
1832 return Err(data_error("Dataset.set_attrs: attrs must be a struct"));
1833 };
1834 let root = dataset_root(&path);
1835 let mut manifest = read_manifest_async(&root).await?;
1836 for (k, v) in incoming.fields {
1837 manifest.attrs.insert(k, value_to_json(&v));
1838 }
1839 manifest.updated_at = now_rfc3339();
1840 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
1841 write_manifest_async(&root, &manifest).await?;
1842 Ok(Value::Bool(true))
1843}
1844
1845#[runtime_builtin(
1846 name = "Dataset.begin",
1847 category = "io/data",
1848 type_resolver(crate::builtins::io::type_resolvers::data_tx_type),
1849 descriptor(crate::builtins::io::data::DATASET_BEGIN_DESCRIPTOR),
1850 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1851 builtin_path = "crate::builtins::io::data"
1852)]
1853async fn dataset_begin_builtin(base: Value, _rest: Vec<Value>) -> BuiltinResult<Value> {
1854 let path = dataset_path_from_object(&base, "Dataset.begin")?;
1855 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1856 let tx_id = start_tx(path.clone(), manifest.txn_sequence)?;
1857 tracing::info!(
1858 target: "runmat.data",
1859 dataset = path,
1860 tx_id = tx_id,
1861 base_sequence = manifest.txn_sequence,
1862 "data transaction begin"
1863 );
1864 Ok(transaction_object(&path, &tx_id))
1865}
1866
1867#[runtime_builtin(
1868 name = "Dataset.snapshot",
1869 category = "io/data",
1870 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1871 descriptor(crate::builtins::io::data::DATASET_SNAPSHOT_DESCRIPTOR),
1872 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1873 builtin_path = "crate::builtins::io::data"
1874)]
1875async fn dataset_snapshot_builtin(
1876 base: Value,
1877 label: Value,
1878 _rest: Vec<Value>,
1879) -> BuiltinResult<Value> {
1880 let path = dataset_path_from_object(&base, "Dataset.snapshot")?;
1881 let label = parse_string(&label, "Dataset.snapshot label")?;
1882 let root = dataset_root(&path);
1883 let snapshots = root.join(".snapshots");
1884 runmat_filesystem::create_dir_all_async(&snapshots)
1885 .await
1886 .map_err(|err| {
1887 data_error(format!(
1888 "Dataset.snapshot: failed to create snapshots dir: {err}"
1889 ))
1890 })?;
1891 let src = manifest_path(&root);
1892 let dst = snapshots.join(format!("{}.manifest.json", sanitize_label(&label)));
1893 copy_file(&src, &dst).await?;
1894 Ok(Value::String(dst.to_string_lossy().to_string()))
1895}
1896
1897#[runtime_builtin(
1898 name = "Dataset.refresh",
1899 category = "io/data",
1900 type_resolver(crate::builtins::io::type_resolvers::data_dataset_type),
1901 descriptor(crate::builtins::io::data::DATASET_REFRESH_DESCRIPTOR),
1902 integer_audit(crate::builtins::io::data::DATASET_NONATTRIBUTE_INTEGER_AUDIT),
1903 builtin_path = "crate::builtins::io::data"
1904)]
1905async fn dataset_refresh_builtin(base: Value) -> BuiltinResult<Value> {
1906 let path = dataset_path_from_object(&base, "Dataset.refresh")?;
1907 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1908 let mut ds = dataset_object(&path, &manifest);
1909 hydrate_dataset_descriptor_async(&path, &mut ds).await;
1910 Ok(ds)
1911}
1912
1913#[runtime_builtin(
1914 name = "DataArray.name",
1915 category = "io/data",
1916 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1917 descriptor(crate::builtins::io::data::DATAARRAY_NAME_DESCRIPTOR),
1918 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
1919 builtin_path = "crate::builtins::io::data"
1920)]
1921async fn data_array_name_builtin(base: Value) -> BuiltinResult<Value> {
1922 let obj = as_object(&base, "DataArray.name")?;
1923 Ok(get_object_prop(obj, "__array_name")?.clone())
1924}
1925
1926#[runtime_builtin(
1927 name = "DataArray.dtype",
1928 category = "io/data",
1929 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
1930 descriptor(crate::builtins::io::data::DATAARRAY_DTYPE_DESCRIPTOR),
1931 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
1932 builtin_path = "crate::builtins::io::data"
1933)]
1934async fn data_array_dtype_builtin(base: Value) -> BuiltinResult<Value> {
1935 let (path, name) = array_identity(&base, "DataArray.dtype")?;
1936 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1937 let meta = manifest
1938 .arrays
1939 .get(&name)
1940 .ok_or_else(|| data_error(format!("DataArray.dtype: array '{name}' not found")))?;
1941 Ok(Value::String(meta.dtype.clone()))
1942}
1943
1944#[runtime_builtin(
1945 name = "DataArray.shape",
1946 category = "io/data",
1947 type_resolver(crate::builtins::io::type_resolvers::data_shape_tensor_type),
1948 descriptor(crate::builtins::io::data::DATAARRAY_SHAPE_DESCRIPTOR),
1949 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
1950 builtin_path = "crate::builtins::io::data"
1951)]
1952async fn data_array_shape_builtin(base: Value) -> BuiltinResult<Value> {
1953 let (path, name) = array_identity(&base, "DataArray.shape")?;
1954 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1955 let meta = manifest
1956 .arrays
1957 .get(&name)
1958 .ok_or_else(|| data_error(format!("DataArray.shape: array '{name}' not found")))?;
1959 let values = meta.shape.iter().map(|v| *v as f64).collect::<Vec<_>>();
1960 let tensor = Tensor::new(values, vec![1, meta.shape.len()])
1961 .map_err(|err| data_error(format!("DataArray.shape: {err}")))?;
1962 Ok(Value::Tensor(tensor))
1963}
1964
1965#[runtime_builtin(
1966 name = "DataArray.rank",
1967 category = "io/data",
1968 type_resolver(crate::builtins::io::type_resolvers::data_int_type),
1969 descriptor(crate::builtins::io::data::DATAARRAY_RANK_DESCRIPTOR),
1970 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
1971 builtin_path = "crate::builtins::io::data"
1972)]
1973async fn data_array_rank_builtin(base: Value) -> BuiltinResult<Value> {
1974 let (path, name) = array_identity(&base, "DataArray.rank")?;
1975 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1976 let meta = manifest
1977 .arrays
1978 .get(&name)
1979 .ok_or_else(|| data_error(format!("DataArray.rank: array '{name}' not found")))?;
1980 Ok(Value::Num(meta.shape.len() as f64))
1981}
1982
1983#[runtime_builtin(
1984 name = "DataArray.chunk_shape",
1985 category = "io/data",
1986 type_resolver(crate::builtins::io::type_resolvers::data_shape_tensor_type),
1987 descriptor(crate::builtins::io::data::DATAARRAY_CHUNK_SHAPE_DESCRIPTOR),
1988 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
1989 builtin_path = "crate::builtins::io::data"
1990)]
1991async fn data_array_chunk_shape_builtin(base: Value) -> BuiltinResult<Value> {
1992 let (path, name) = array_identity(&base, "DataArray.chunk_shape")?;
1993 let manifest = read_manifest_async(&dataset_root(&path)).await?;
1994 let meta = manifest
1995 .arrays
1996 .get(&name)
1997 .ok_or_else(|| data_error(format!("DataArray.chunk_shape: array '{name}' not found")))?;
1998 let values = meta
1999 .chunk_shape
2000 .iter()
2001 .map(|v| *v as f64)
2002 .collect::<Vec<_>>();
2003 let tensor = Tensor::new(values, vec![1, meta.chunk_shape.len()])
2004 .map_err(|err| data_error(format!("DataArray.chunk_shape: {err}")))?;
2005 Ok(Value::Tensor(tensor))
2006}
2007
2008#[runtime_builtin(
2009 name = "DataArray.codec",
2010 category = "io/data",
2011 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
2012 descriptor(crate::builtins::io::data::DATAARRAY_CODEC_DESCRIPTOR),
2013 integer_audit(crate::builtins::io::data::DATAARRAY_METADATA_INTEGER_AUDIT),
2014 builtin_path = "crate::builtins::io::data"
2015)]
2016async fn data_array_codec_builtin(base: Value) -> BuiltinResult<Value> {
2017 let (path, name) = array_identity(&base, "DataArray.codec")?;
2018 let manifest = read_manifest_async(&dataset_root(&path)).await?;
2019 let meta = manifest
2020 .arrays
2021 .get(&name)
2022 .ok_or_else(|| data_error(format!("DataArray.codec: array '{name}' not found")))?;
2023 Ok(Value::String(meta.codec.clone()))
2024}
2025
2026#[runtime_builtin(
2027 name = "DataArray.read",
2028 category = "io/data",
2029 type_resolver(crate::builtins::io::type_resolvers::data_tensor_type),
2030 descriptor(crate::builtins::io::data::DATAARRAY_READ_DESCRIPTOR),
2031 integer_capabilities(crate::builtins::io::data::DATAARRAY_READ_INTEGER_CAPABILITIES),
2032 builtin_path = "crate::builtins::io::data"
2033)]
2034async fn data_array_read_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
2035 let (path, name) = array_identity(&base, "DataArray.read")?;
2036 let root = dataset_root(&path);
2037 let manifest = read_manifest_async(&root).await?;
2038 let meta = manifest
2039 .arrays
2040 .get(&name)
2041 .ok_or_else(|| data_error(format!("DataArray.read: array '{name}' not found")))?;
2042 let payload = read_array_payload_async(&root, meta).await?;
2043 let sliced = if let Some(slice_spec) = rest.first() {
2044 read_slice_payload(&payload, slice_spec)?
2045 } else {
2046 payload
2047 };
2048 sliced
2049 .into_value()
2050 .map_err(|err| data_error(format!("DataArray.read: {err}")))
2051}
2052
2053#[runtime_builtin(
2054 name = "DataArray.write",
2055 category = "io/data",
2056 sink = true,
2057 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2058 descriptor(crate::builtins::io::data::DATAARRAY_WRITE_DESCRIPTOR),
2059 integer_capabilities(crate::builtins::io::data::DATAARRAY_WRITE_INTEGER_CAPABILITIES),
2060 builtin_path = "crate::builtins::io::data"
2061)]
2062async fn data_array_write_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
2063 let (path, name) = array_identity(&base, "DataArray.write")?;
2064 let (slice_spec, value) = match rest.as_slice() {
2065 [v] => (None, v),
2066 [slice, v] => (Some(slice), v),
2067 _ => {
2068 return Err(data_error(
2069 "DataArray.write expects values or (sliceSpec, values) arguments",
2070 ))
2071 }
2072 };
2073 write_array_full_async(&path, &name, slice_spec, value).await?;
2074 Ok(Value::Bool(true))
2075}
2076
2077#[runtime_builtin(
2078 name = "DataArray.resize",
2079 category = "io/data",
2080 sink = true,
2081 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2082 descriptor(crate::builtins::io::data::DATAARRAY_RESIZE_DESCRIPTOR),
2083 integer_capabilities(crate::builtins::io::data::DATAARRAY_RESIZE_INTEGER_CAPABILITIES),
2084 builtin_path = "crate::builtins::io::data"
2085)]
2086async fn data_array_resize_builtin(
2087 base: Value,
2088 new_shape: Value,
2089 _rest: Vec<Value>,
2090) -> BuiltinResult<Value> {
2091 let (path, name) = array_identity(&base, "DataArray.resize")?;
2092 let shape = parse_shape_from_value(&new_shape)?;
2093 let root = dataset_root(&path);
2094 let mut manifest = read_manifest_async(&root).await?;
2095 let meta = manifest
2096 .arrays
2097 .get_mut(&name)
2098 .ok_or_else(|| data_error(format!("DataArray.resize: array '{name}' not found")))?;
2099 meta.shape = shape.clone();
2100 let payload = DataArrayPayload::zeros(meta.dtype.clone(), shape.clone())?;
2101 let (payload_path, chunk_index_path) =
2102 write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
2103 meta.data_path = make_rel_data_path(&root, &payload_path)?;
2104 meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
2105 manifest.updated_at = now_rfc3339();
2106 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2107 write_manifest_async(&root, &manifest).await?;
2108 Ok(Value::Bool(true))
2109}
2110
2111#[runtime_builtin(
2112 name = "DataArray.fill",
2113 category = "io/data",
2114 sink = true,
2115 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2116 descriptor(crate::builtins::io::data::DATAARRAY_FILL_DESCRIPTOR),
2117 integer_capabilities(crate::builtins::io::data::DATAARRAY_FILL_INTEGER_CAPABILITIES),
2118 builtin_path = "crate::builtins::io::data"
2119)]
2120async fn data_array_fill_builtin(
2121 base: Value,
2122 value: Value,
2123 _rest: Vec<Value>,
2124) -> BuiltinResult<Value> {
2125 let (path, name) = array_identity(&base, "DataArray.fill")?;
2126 let root = dataset_root(&path);
2127 let mut manifest = read_manifest_async(&root).await?;
2128 let meta = manifest
2129 .arrays
2130 .get_mut(&name)
2131 .ok_or_else(|| data_error(format!("DataArray.fill: array '{name}' not found")))?;
2132 let payload = DataArrayPayload::filled(meta.dtype.clone(), meta.shape.clone(), &value)?;
2133 let (payload_path, chunk_index_path) =
2134 write_array_payload_async(&root, &name, &payload, &meta.chunk_shape).await?;
2135 meta.data_path = make_rel_data_path(&root, &payload_path)?;
2136 meta.chunk_index_path = Some(make_rel_data_path(&root, &chunk_index_path)?);
2137 manifest.updated_at = now_rfc3339();
2138 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2139 write_manifest_async(&root, &manifest).await?;
2140 Ok(Value::Bool(true))
2141}
2142
2143#[runtime_builtin(
2144 name = "DataTransaction.id",
2145 category = "io/data",
2146 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
2147 descriptor(crate::builtins::io::data::DATATX_ID_DESCRIPTOR),
2148 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2149 builtin_path = "crate::builtins::io::data"
2150)]
2151async fn data_tx_id_builtin(base: Value) -> BuiltinResult<Value> {
2152 let obj = as_object(&base, "DataTransaction.id")?;
2153 Ok(get_object_prop(obj, "__tx_id")?.clone())
2154}
2155
2156#[runtime_builtin(
2157 name = "DataTransaction.write",
2158 category = "io/data",
2159 sink = true,
2160 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2161 descriptor(crate::builtins::io::data::DATATX_WRITE_DESCRIPTOR),
2162 integer_capabilities(crate::builtins::io::data::DATATX_WRITE_INTEGER_CAPABILITIES),
2163 builtin_path = "crate::builtins::io::data"
2164)]
2165async fn data_tx_write_builtin(
2166 base: Value,
2167 array_name: Value,
2168 slice: Value,
2169 values: Value,
2170 _rest: Vec<Value>,
2171) -> BuiltinResult<Value> {
2172 let tx_id = tx_id_from_object(&base, "DataTransaction.write")?;
2173 let array_name = parse_string(&array_name, "DataTransaction.write arrayName")?;
2174 with_tx_mut(&tx_id, |tx| {
2175 if tx.status != TxnStatus::Open {
2176 return Err(data_error("DataTransaction.write: transaction is not open"));
2177 }
2178 tx.writes.push(PendingWrite {
2179 array: array_name,
2180 slice_spec: Some(slice),
2181 value: values,
2182 });
2183 Ok(())
2184 })?;
2185 Ok(Value::Bool(true))
2186}
2187
2188#[runtime_builtin(
2189 name = "DataTransaction.set_attr",
2190 category = "io/data",
2191 sink = true,
2192 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2193 descriptor(crate::builtins::io::data::DATATX_SET_ATTR_DESCRIPTOR),
2194 integer_capabilities(crate::builtins::io::data::DATATX_SET_ATTR_INTEGER_CAPABILITIES),
2195 builtin_path = "crate::builtins::io::data"
2196)]
2197async fn data_tx_set_attr_builtin(base: Value, key: Value, value: Value) -> BuiltinResult<Value> {
2198 let tx_id = tx_id_from_object(&base, "DataTransaction.set_attr")?;
2199 let key = parse_string(&key, "DataTransaction.set_attr key")?;
2200 with_tx_mut(&tx_id, |tx| {
2201 if tx.status != TxnStatus::Open {
2202 return Err(data_error(
2203 "DataTransaction.set_attr: transaction is not open",
2204 ));
2205 }
2206 tx.attrs.insert(key, value);
2207 Ok(())
2208 })?;
2209 Ok(Value::Bool(true))
2210}
2211
2212#[runtime_builtin(
2213 name = "DataTransaction.set_attrs",
2214 category = "io/data",
2215 sink = true,
2216 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2217 descriptor(crate::builtins::io::data::DATATX_SET_ATTRS_DESCRIPTOR),
2218 integer_capabilities(crate::builtins::io::data::DATATX_SET_ATTRS_INTEGER_CAPABILITIES),
2219 builtin_path = "crate::builtins::io::data"
2220)]
2221async fn data_tx_set_attrs_builtin(base: Value, attrs: Value) -> BuiltinResult<Value> {
2222 let tx_id = tx_id_from_object(&base, "DataTransaction.set_attrs")?;
2223 let Value::Struct(incoming) = attrs else {
2224 return Err(data_error(
2225 "DataTransaction.set_attrs: attrs must be struct",
2226 ));
2227 };
2228 with_tx_mut(&tx_id, |tx| {
2229 if tx.status != TxnStatus::Open {
2230 return Err(data_error(
2231 "DataTransaction.set_attrs: transaction is not open",
2232 ));
2233 }
2234 for (k, v) in incoming.fields {
2235 tx.attrs.insert(k, v);
2236 }
2237 Ok(())
2238 })?;
2239 Ok(Value::Bool(true))
2240}
2241
2242#[runtime_builtin(
2243 name = "DataTransaction.resize",
2244 category = "io/data",
2245 sink = true,
2246 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2247 descriptor(crate::builtins::io::data::DATATX_RESIZE_DESCRIPTOR),
2248 integer_capabilities(crate::builtins::io::data::DATATX_RESIZE_INTEGER_CAPABILITIES),
2249 builtin_path = "crate::builtins::io::data"
2250)]
2251async fn data_tx_resize_builtin(
2252 base: Value,
2253 array_name: Value,
2254 new_shape: Value,
2255 _rest: Vec<Value>,
2256) -> BuiltinResult<Value> {
2257 let tx_id = tx_id_from_object(&base, "DataTransaction.resize")?;
2258 let array_name = parse_string(&array_name, "DataTransaction.resize arrayName")?;
2259 let shape = parse_shape_from_value(&new_shape)?;
2260 with_tx_mut(&tx_id, |tx| {
2261 if tx.status != TxnStatus::Open {
2262 return Err(data_error(
2263 "DataTransaction.resize: transaction is not open",
2264 ));
2265 }
2266 tx.resizes.push(PendingResize {
2267 array: array_name,
2268 shape,
2269 });
2270 Ok(())
2271 })?;
2272 Ok(Value::Bool(true))
2273}
2274
2275#[runtime_builtin(
2276 name = "DataTransaction.fill",
2277 category = "io/data",
2278 sink = true,
2279 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2280 descriptor(crate::builtins::io::data::DATATX_FILL_DESCRIPTOR),
2281 integer_capabilities(crate::builtins::io::data::DATATX_FILL_INTEGER_CAPABILITIES),
2282 builtin_path = "crate::builtins::io::data"
2283)]
2284async fn data_tx_fill_builtin(
2285 base: Value,
2286 array_name: Value,
2287 value: Value,
2288 rest: Vec<Value>,
2289) -> BuiltinResult<Value> {
2290 let tx_id = tx_id_from_object(&base, "DataTransaction.fill")?;
2291 let array_name = parse_string(&array_name, "DataTransaction.fill arrayName")?;
2292 let slice_spec = rest.first().cloned();
2293 with_tx_mut(&tx_id, |tx| {
2294 if tx.status != TxnStatus::Open {
2295 return Err(data_error("DataTransaction.fill: transaction is not open"));
2296 }
2297 tx.fills.push(PendingFill {
2298 array: array_name,
2299 slice_spec,
2300 value,
2301 });
2302 Ok(())
2303 })?;
2304 Ok(Value::Bool(true))
2305}
2306
2307#[runtime_builtin(
2308 name = "DataTransaction.delete_array",
2309 category = "io/data",
2310 sink = true,
2311 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2312 descriptor(crate::builtins::io::data::DATATX_DELETE_ARRAY_DESCRIPTOR),
2313 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2314 builtin_path = "crate::builtins::io::data"
2315)]
2316async fn data_tx_delete_array_builtin(base: Value, array_name: Value) -> BuiltinResult<Value> {
2317 let tx_id = tx_id_from_object(&base, "DataTransaction.delete_array")?;
2318 let array_name = parse_string(&array_name, "DataTransaction.delete_array arrayName")?;
2319 with_tx_mut(&tx_id, |tx| {
2320 if tx.status != TxnStatus::Open {
2321 return Err(data_error(
2322 "DataTransaction.delete_array: transaction is not open",
2323 ));
2324 }
2325 tx.delete_arrays.push(array_name);
2326 Ok(())
2327 })?;
2328 Ok(Value::Bool(true))
2329}
2330
2331#[runtime_builtin(
2332 name = "DataTransaction.create_array",
2333 category = "io/data",
2334 sink = true,
2335 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2336 descriptor(crate::builtins::io::data::DATATX_CREATE_ARRAY_DESCRIPTOR),
2337 integer_capabilities(crate::builtins::io::data::DATATX_CREATE_ARRAY_INTEGER_CAPABILITIES),
2338 builtin_path = "crate::builtins::io::data"
2339)]
2340async fn data_tx_create_array_builtin(
2341 base: Value,
2342 array_name: Value,
2343 meta: Value,
2344) -> BuiltinResult<Value> {
2345 let tx_id = tx_id_from_object(&base, "DataTransaction.create_array")?;
2346 let array_name = parse_string(&array_name, "DataTransaction.create_array arrayName")?;
2347 let meta = parse_array_meta(&array_name, &meta)?;
2348 with_tx_mut(&tx_id, |tx| {
2349 if tx.status != TxnStatus::Open {
2350 return Err(data_error(
2351 "DataTransaction.create_array: transaction is not open",
2352 ));
2353 }
2354 tx.create_arrays.push(PendingCreateArray {
2355 array: array_name,
2356 meta,
2357 });
2358 Ok(())
2359 })?;
2360 Ok(Value::Bool(true))
2361}
2362
2363#[runtime_builtin(
2364 name = "DataTransaction.commit",
2365 category = "io/data",
2366 sink = true,
2367 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2368 descriptor(crate::builtins::io::data::DATATX_COMMIT_DESCRIPTOR),
2369 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2370 builtin_path = "crate::builtins::io::data"
2371)]
2372async fn data_tx_commit_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
2373 let options = match rest.as_slice() {
2374 [] => None,
2375 [Value::Struct(options)] => Some(options),
2376 [_] => {
2377 return Err(data_error(
2378 "DataTransaction.commit: options must be a scalar struct",
2379 ))
2380 }
2381 _ => {
2382 return Err(data_error(
2383 "DataTransaction.commit: expected at most one options struct",
2384 ))
2385 }
2386 };
2387 let tx_id = tx_id_from_object(&base, "DataTransaction.commit")?;
2388 let (dataset_path, base_sequence, writes, resizes, fills, create_arrays, delete_arrays, attrs) =
2389 with_tx(&tx_id, |tx| {
2390 if tx.status != TxnStatus::Open {
2391 return Err(data_error(
2392 "DataTransaction.commit: transaction is not open",
2393 ));
2394 }
2395 Ok((
2396 tx.dataset_path.clone(),
2397 tx.base_sequence,
2398 tx.writes.clone(),
2399 tx.resizes.clone(),
2400 tx.fills.clone(),
2401 tx.create_arrays.clone(),
2402 tx.delete_arrays.clone(),
2403 tx.attrs.clone(),
2404 ))
2405 })?;
2406
2407 let write_ops = writes.len();
2408 let resize_ops = resizes.len();
2409 let fill_ops = fills.len();
2410 let create_ops = create_arrays.len();
2411 let delete_ops = delete_arrays.len();
2412 let attr_updates = attrs.len();
2413
2414 let root = dataset_root(&dataset_path);
2415 let mut manifest = read_manifest_async(&root).await?;
2416 ensure_manifest_sequence(base_sequence, &manifest)?;
2417 if let Some(options) = options {
2418 if let Some(expected) = options.fields.get("if_manifest") {
2419 let expected = parse_string(expected, "DataTransaction.commit if_manifest")?;
2420 let actual = manifest_version_token(&manifest);
2421 if expected != actual {
2422 tracing::warn!(
2423 target: "runmat.data",
2424 tx_id = tx_id,
2425 expected_manifest = expected,
2426 actual_manifest = actual,
2427 "data transaction manifest conflict"
2428 );
2429 return Err(data_error(
2430 "MANIFEST_CONFLICT: if_manifest precondition failed",
2431 ));
2432 }
2433 }
2434 }
2435 for create in create_arrays {
2436 create_array_in_manifest(&root, &mut manifest, &create.array, create.meta).await?;
2437 }
2438 for resize in resizes {
2439 resize_array_in_manifest(&root, &mut manifest, &resize.array, resize.shape).await?;
2440 }
2441 for fill in fills {
2442 fill_array_in_manifest(
2443 &root,
2444 &mut manifest,
2445 &fill.array,
2446 fill.slice_spec.as_ref(),
2447 &fill.value,
2448 )
2449 .await?;
2450 }
2451 for write in writes {
2452 apply_write_to_manifest_async(
2453 &root,
2454 &mut manifest,
2455 &write.array,
2456 write.slice_spec.as_ref(),
2457 &write.value,
2458 )
2459 .await?;
2460 }
2461 for array_name in delete_arrays {
2462 delete_array_in_manifest_async(&root, &mut manifest, &array_name).await?;
2463 }
2464 for (k, v) in attrs {
2465 manifest.attrs.insert(k, value_to_json(&v));
2466 }
2467 manifest.updated_at = now_rfc3339();
2468 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2469 write_manifest_async(&root, &manifest).await?;
2470 with_tx_mut(&tx_id, |tx| {
2471 tx.status = TxnStatus::Committed;
2472 Ok(())
2473 })?;
2474 tracing::info!(
2475 target: "runmat.data",
2476 dataset = dataset_path,
2477 tx_id = tx_id,
2478 write_ops = write_ops,
2479 resize_ops = resize_ops,
2480 fill_ops = fill_ops,
2481 create_ops = create_ops,
2482 delete_ops = delete_ops,
2483 attr_updates = attr_updates,
2484 next_sequence = manifest.txn_sequence,
2485 "data transaction commit"
2486 );
2487 remove_tx(&tx_id)?;
2488 Ok(Value::Bool(true))
2489}
2490
2491#[runtime_builtin(
2492 name = "commit",
2493 category = "io/data",
2494 sink = true,
2495 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2496 descriptor(crate::builtins::io::data::COMMIT_ALIAS_DESCRIPTOR),
2497 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2498 builtin_path = "crate::builtins::io::data"
2499)]
2500async fn data_tx_commit_alias_builtin(base: Value, rest: Vec<Value>) -> BuiltinResult<Value> {
2501 match &base {
2502 Value::Object(obj) if obj.class_name == "DataTransaction" => {
2503 data_tx_commit_builtin(base, rest).await
2504 }
2505 _ => Err(data_error(
2506 "commit: receiver must be a DataTransaction (use tx = ds.begin())",
2507 )),
2508 }
2509}
2510
2511#[runtime_builtin(
2512 name = "DataTransaction.abort",
2513 category = "io/data",
2514 sink = true,
2515 type_resolver(crate::builtins::io::type_resolvers::data_bool_type),
2516 descriptor(crate::builtins::io::data::DATATX_ABORT_DESCRIPTOR),
2517 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2518 builtin_path = "crate::builtins::io::data"
2519)]
2520async fn data_tx_abort_builtin(base: Value) -> BuiltinResult<Value> {
2521 let tx_id = tx_id_from_object(&base, "DataTransaction.abort")?;
2522 with_tx_mut(&tx_id, |tx| {
2523 tx.status = TxnStatus::Aborted;
2524 Ok(())
2525 })?;
2526 tracing::info!(
2527 target: "runmat.data",
2528 tx_id = tx_id,
2529 "data transaction abort"
2530 );
2531 remove_tx(&tx_id)?;
2532 Ok(Value::Bool(true))
2533}
2534
2535#[runtime_builtin(
2536 name = "DataTransaction.status",
2537 category = "io/data",
2538 type_resolver(crate::builtins::io::type_resolvers::data_string_type),
2539 descriptor(crate::builtins::io::data::DATATX_STATUS_DESCRIPTOR),
2540 integer_audit(crate::builtins::io::data::DATATX_LIFECYCLE_INTEGER_AUDIT),
2541 builtin_path = "crate::builtins::io::data"
2542)]
2543async fn data_tx_status_builtin(base: Value) -> BuiltinResult<Value> {
2544 let tx_id = tx_id_from_object(&base, "DataTransaction.status")?;
2545 with_tx(&tx_id, |tx| {
2546 let status = match tx.status {
2547 TxnStatus::Open => "open",
2548 TxnStatus::Committed => "committed",
2549 TxnStatus::Aborted => "aborted",
2550 };
2551 Ok(Value::String(status.to_string()))
2552 })
2553}
2554
2555fn dataset_path_from_object(base: &Value, context: &str) -> BuiltinResult<String> {
2556 let obj = as_object(base, context)?;
2557 parse_string(get_object_prop(obj, "__data_path")?, context)
2558}
2559
2560fn tx_id_from_object(base: &Value, context: &str) -> BuiltinResult<String> {
2561 let obj = as_object(base, context)?;
2562 parse_string(get_object_prop(obj, "__tx_id")?, context)
2563}
2564
2565fn array_identity(base: &Value, context: &str) -> BuiltinResult<(String, String)> {
2566 let obj = as_object(base, context)?;
2567 let path = parse_string(get_object_prop(obj, "__data_path")?, context)?;
2568 let name = parse_string(get_object_prop(obj, "__array_name")?, context)?;
2569 Ok((path, name))
2570}
2571
2572fn as_object<'a>(value: &'a Value, context: &str) -> BuiltinResult<&'a ObjectInstance> {
2573 match value {
2574 Value::Object(obj) => Ok(obj),
2575 _ => Err(data_error(format!("{context}: expected object receiver"))),
2576 }
2577}
2578
2579async fn hydrate_dataset_descriptor_async(path: &str, dataset: &mut Value) {
2580 let request = runmat_filesystem::data_contract::DataManifestRequest {
2581 path: path.to_string(),
2582 version: None,
2583 };
2584 let descriptor = match runmat_filesystem::data_manifest_descriptor_async(&request).await {
2585 Ok(descriptor) => descriptor,
2586 Err(_) => return,
2587 };
2588 let Value::Object(obj) = dataset else {
2589 return;
2590 };
2591 if !descriptor.dataset_id.is_empty() {
2592 obj.properties.insert(
2593 "__data_id".to_string(),
2594 Value::String(descriptor.dataset_id),
2595 );
2596 }
2597 obj.properties.insert(
2598 "__data_version".to_string(),
2599 Value::String(format!(
2600 "{}:{}",
2601 descriptor.updated_at, descriptor.txn_sequence
2602 )),
2603 );
2604}
2605
2606fn sanitize_label(label: &str) -> String {
2607 label
2608 .chars()
2609 .map(|ch| {
2610 if ch.is_ascii_alphanumeric() || ch == '-' || ch == '_' {
2611 ch
2612 } else {
2613 '_'
2614 }
2615 })
2616 .collect()
2617}
2618
2619async fn copy_file(src: &PathBuf, dst: &PathBuf) -> BuiltinResult<()> {
2620 let bytes = runmat_filesystem::read_async(src)
2621 .await
2622 .map_err(|err| data_error(format!("failed to open '{}': {err}", src.display())))?;
2623 let parent = dst.parent().ok_or_else(|| {
2624 data_error(format!(
2625 "invalid destination path '{}': missing parent",
2626 dst.display()
2627 ))
2628 })?;
2629 runmat_filesystem::create_dir_all_async(parent)
2630 .await
2631 .map_err(|err| data_error(format!("failed to create '{}': {err}", parent.display())))?;
2632 runmat_filesystem::write_async(dst, &bytes)
2633 .await
2634 .map_err(|err| {
2635 data_error(format!(
2636 "failed to copy '{}' -> '{}': {err}",
2637 src.display(),
2638 dst.display()
2639 ))
2640 })?;
2641 Ok(())
2642}
2643
2644fn make_rel_data_path(
2645 root: &std::path::Path,
2646 payload_path: &std::path::Path,
2647) -> BuiltinResult<String> {
2648 let rel = payload_path
2649 .strip_prefix(root)
2650 .map_err(|err| data_error(format!("failed to compute relative data path: {err}")))?;
2651 Ok(rel.to_string_lossy().to_string())
2652}
2653
2654async fn create_array_in_manifest(
2655 root: &std::path::Path,
2656 manifest: &mut DataManifest,
2657 array_name: &str,
2658 mut meta: DataArrayMeta,
2659) -> BuiltinResult<()> {
2660 if manifest.arrays.contains_key(array_name) {
2661 return Err(data_error(format!(
2662 "DataTransaction.create_array: array '{array_name}' already exists"
2663 )));
2664 }
2665 let payload = DataArrayPayload::zeros(meta.dtype.clone(), meta.shape.clone())?;
2666 let (payload_path, chunk_index_path) =
2667 write_array_payload_async(root, array_name, &payload, &meta.chunk_shape).await?;
2668 meta.data_path = make_rel_data_path(root, &payload_path)?;
2669 meta.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2670 manifest.arrays.insert(array_name.to_string(), meta);
2671 Ok(())
2672}
2673
2674async fn resize_array_in_manifest(
2675 root: &std::path::Path,
2676 manifest: &mut DataManifest,
2677 array_name: &str,
2678 shape: Vec<usize>,
2679) -> BuiltinResult<()> {
2680 let meta = manifest
2681 .arrays
2682 .get_mut(array_name)
2683 .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2684 meta.shape = shape.clone();
2685 let payload = DataArrayPayload::zeros(meta.dtype.clone(), shape.clone())?;
2686 let (payload_path, chunk_index_path) =
2687 write_array_payload_async(root, array_name, &payload, &meta.chunk_shape).await?;
2688 meta.data_path = make_rel_data_path(root, &payload_path)?;
2689 meta.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2690 Ok(())
2691}
2692
2693async fn fill_array_in_manifest(
2694 root: &std::path::Path,
2695 manifest: &mut DataManifest,
2696 array_name: &str,
2697 slice_spec: Option<&Value>,
2698 value: &Value,
2699) -> BuiltinResult<()> {
2700 let meta: DataArrayMeta = manifest
2701 .arrays
2702 .get(array_name)
2703 .cloned()
2704 .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2705 let payload = read_array_payload_async(root, &meta).await?;
2706 let next_payload = if let Some(slice_spec) = slice_spec {
2707 let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
2708 let target_shape: Vec<usize> = ranges
2709 .iter()
2710 .map(|r| r.end.saturating_sub(r.start))
2711 .collect();
2712 let rhs = DataArrayPayload::filled(payload.dtype.clone(), target_shape, value)?
2713 .into_value()
2714 .map_err(|err| data_error(format!("DataTransaction.fill: {err}")))?;
2715 write_slice_payload(&payload, slice_spec, &rhs)?
2716 } else {
2717 DataArrayPayload::filled(payload.dtype.clone(), payload.shape.clone(), value)?
2718 };
2719 let (payload_path, chunk_index_path) =
2720 write_array_payload_async(root, array_name, &next_payload, &meta.chunk_shape).await?;
2721 if let Some(updated) = manifest.arrays.get_mut(array_name) {
2722 updated.shape = next_payload.shape.clone();
2723 updated.data_path = make_rel_data_path(root, &payload_path)?;
2724 updated.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2725 }
2726 Ok(())
2727}
2728
2729async fn delete_array_in_manifest_async(
2730 root: &std::path::Path,
2731 manifest: &mut DataManifest,
2732 array_name: &str,
2733) -> BuiltinResult<()> {
2734 let removed = manifest.arrays.remove(array_name);
2735 if removed.is_none() {
2736 return Err(data_error(format!(
2737 "DataTransaction.delete_array: array '{array_name}' not found"
2738 )));
2739 }
2740 let array_dir = root.join("arrays").join(array_name);
2741 if runmat_filesystem::metadata_async(&array_dir).await.is_ok() {
2742 runmat_filesystem::remove_dir_all_async(&array_dir)
2743 .await
2744 .map_err(|err| {
2745 data_error(format!(
2746 "DataTransaction.delete_array: failed to remove '{}': {err}",
2747 array_dir.display()
2748 ))
2749 })?;
2750 }
2751 Ok(())
2752}
2753
2754fn parse_array_meta(array_name: &str, meta: &Value) -> BuiltinResult<DataArrayMeta> {
2755 let Value::Struct(meta_struct) = meta else {
2756 return Err(data_error(
2757 "DataTransaction.create_array: meta must be a struct",
2758 ));
2759 };
2760 let dtype = meta_struct
2761 .fields
2762 .get("dtype")
2763 .map(|v| parse_string(v, "DataTransaction.create_array dtype"))
2764 .transpose()?
2765 .unwrap_or_else(|| "f64".to_string());
2766 let shape = meta_struct
2767 .fields
2768 .get("shape")
2769 .map(parse_shape_from_value)
2770 .transpose()?
2771 .unwrap_or_else(|| vec![0, 0]);
2772 let chunk_shape = meta_struct
2773 .fields
2774 .get("chunk")
2775 .map(parse_shape_from_value)
2776 .transpose()?
2777 .unwrap_or_else(|| default_chunk_shape(&shape));
2778 validate_chunk_shape(&shape, &chunk_shape)?;
2779 let codec = meta_struct
2780 .fields
2781 .get("codec")
2782 .map(|v| parse_string(v, "DataTransaction.create_array codec"))
2783 .transpose()?
2784 .unwrap_or_else(|| "zstd".to_string());
2785 Ok(DataArrayMeta {
2786 dtype,
2787 shape,
2788 chunk_shape,
2789 order: "column_major".to_string(),
2790 codec,
2791 chunk_index_path: Some(format!("arrays/{array_name}/chunks/index.json")),
2792 data_path: format!("arrays/{array_name}/data.f64.json"),
2793 })
2794}
2795
2796fn default_chunk_shape(shape: &[usize]) -> Vec<usize> {
2797 if shape.is_empty() {
2798 return Vec::new();
2799 }
2800 let mut out = shape.to_vec();
2801 if out.len() == 1 {
2802 out[0] = out[0].clamp(1, 65_536);
2803 return out;
2804 }
2805 out[0] = out[0].clamp(1, 256);
2806 out[1] = out[1].clamp(1, 256);
2807 for dim in out.iter_mut().skip(2) {
2808 *dim = (*dim).clamp(1, 8);
2809 }
2810 out
2811}
2812
2813#[async_recursion::async_recursion(?Send)]
2814async fn copy_dir_recursive(src: &PathBuf, dst: &PathBuf) -> BuiltinResult<()> {
2815 let metadata = runmat_filesystem::metadata_async(src)
2816 .await
2817 .map_err(|err| data_error(format!("failed to stat '{}': {err}", src.display())))?;
2818 if !metadata.is_dir() {
2819 return Err(data_error(format!(
2820 "expected dataset directory at '{}'",
2821 src.display()
2822 )));
2823 }
2824 runmat_filesystem::create_dir_all_async(dst)
2825 .await
2826 .map_err(|err| data_error(format!("failed to create '{}': {err}", dst.display())))?;
2827 for entry in runmat_filesystem::read_dir_async(src)
2828 .await
2829 .map_err(|err| data_error(format!("failed to read '{}': {err}", src.display())))?
2830 {
2831 let entry_src = entry.path().to_path_buf();
2832 let entry_dst = dst.join(entry.file_name());
2833 if entry.is_dir() {
2834 copy_dir_recursive(&entry_src, &entry_dst).await?;
2835 continue;
2836 }
2837 copy_file(&entry_src, &entry_dst).await?;
2838 }
2839 Ok(())
2840}
2841
2842fn parse_shape_from_value(value: &Value) -> BuiltinResult<Vec<usize>> {
2843 match value {
2844 Value::Tensor(t) => {
2845 let len = tensor_utils::tensor_element_len(t);
2846 let mut out = Vec::with_capacity(len);
2847 for index in 0..len {
2848 let value = t
2849 .numeric_value_at(index)
2850 .ok_or_else(|| data_error("shape dimensions require valid numeric storage"))?;
2851 out.push(numeric_dimension_to_usize(value).ok_or_else(|| {
2852 data_error(
2853 "shape dimensions must be non-negative finite integers within platform limits",
2854 )
2855 })?);
2856 }
2857 Ok(out)
2858 }
2859 Value::Num(v) => floating_dimension_to_usize(*v)
2860 .map(|value| vec![value])
2861 .ok_or_else(|| {
2862 data_error(
2863 "shape dimensions must be non-negative finite integers within platform limits",
2864 )
2865 }),
2866 Value::Int(v) => v
2867 .try_to_usize()
2868 .map(|n| vec![n])
2869 .ok_or_else(|| data_error("shape dimensions must be non-negative integers")),
2870 _ => Err(data_error("shape must be a numeric vector")),
2871 }
2872}
2873
2874fn floating_dimension_to_usize(value: f64) -> Option<usize> {
2875 if !value.is_finite() || value < 0.0 || value.fract() != 0.0 {
2876 return None;
2877 }
2878 let integer = value as u128;
2879 (integer <= usize::MAX as u128 && integer as f64 == value)
2880 .then(|| usize::try_from(integer).ok())
2881 .flatten()
2882}
2883
2884fn numeric_dimension_to_usize(value: NumericScalar) -> Option<usize> {
2885 match value {
2886 NumericScalar::F64(value) => floating_dimension_to_usize(value),
2887 NumericScalar::F32(value) => floating_dimension_to_usize(f64::from(value)),
2888 value => value
2889 .into_int_value()
2890 .and_then(|value| value.try_to_usize()),
2891 }
2892}
2893
2894async fn write_array_full_async(
2895 dataset_path: &str,
2896 array_name: &str,
2897 slice_spec: Option<&Value>,
2898 value: &Value,
2899) -> BuiltinResult<()> {
2900 let root = dataset_root(dataset_path);
2901 let mut manifest = read_manifest_async(&root).await?;
2902 apply_write_to_manifest_async(&root, &mut manifest, array_name, slice_spec, value).await?;
2903 manifest.updated_at = now_rfc3339();
2904 manifest.txn_sequence = manifest.txn_sequence.saturating_add(1);
2905 write_manifest_async(&root, &manifest).await
2906}
2907
2908async fn apply_write_to_manifest_async(
2909 root: &std::path::Path,
2910 manifest: &mut DataManifest,
2911 array_name: &str,
2912 slice_spec: Option<&Value>,
2913 value: &Value,
2914) -> BuiltinResult<()> {
2915 let meta: DataArrayMeta = manifest
2916 .arrays
2917 .get(array_name)
2918 .cloned()
2919 .ok_or_else(|| data_error(format!("array '{array_name}' not found")))?;
2920
2921 if let Some(slice_spec) = slice_spec {
2922 if apply_slice_write_chunked_async(root, manifest, array_name, &meta, slice_spec, value)
2923 .await?
2924 {
2925 return Ok(());
2926 }
2927 }
2928
2929 let payload = read_array_payload_async(root, &meta).await?;
2930 let next_payload = if let Some(slice_spec) = slice_spec {
2931 write_slice_payload(&payload, slice_spec, value)?
2932 } else {
2933 DataArrayPayload::from_value(payload.dtype.clone(), value)?
2934 };
2935
2936 let (payload_path, chunk_index_path) =
2937 write_array_payload_async(root, array_name, &next_payload, &meta.chunk_shape).await?;
2938 if let Some(updated) = manifest.arrays.get_mut(array_name) {
2939 updated.shape = next_payload.shape.clone();
2940 updated.data_path = make_rel_data_path(root, &payload_path)?;
2941 updated.chunk_index_path = Some(make_rel_data_path(root, &chunk_index_path)?);
2942 }
2943 Ok(())
2944}
2945
2946async fn apply_slice_write_chunked_async(
2947 root: &std::path::Path,
2948 manifest: &mut DataManifest,
2949 array_name: &str,
2950 meta: &DataArrayMeta,
2951 slice_spec: &Value,
2952 value: &Value,
2953) -> BuiltinResult<bool> {
2954 let Some(index_rel_path) = &meta.chunk_index_path else {
2955 return Ok(false);
2956 };
2957 let index_path = root.join(index_rel_path);
2958 if runmat_filesystem::metadata_async(&index_path)
2959 .await
2960 .is_err()
2961 {
2962 return Ok(false);
2963 }
2964 let ranges = parse_slice_spec(slice_spec, &meta.shape)?;
2965 let rhs_shape: Vec<usize> = ranges
2966 .iter()
2967 .map(|r| r.end.saturating_sub(r.start))
2968 .collect();
2969 let rhs = DataArrayPayload::from_value(meta.dtype.clone(), value)?;
2970 if rhs.shape != rhs_shape {
2971 return Err(data_error(format!(
2972 "SHAPE_MISMATCH: rhs shape {:?} must match target slice shape {:?}",
2973 rhs.shape, rhs_shape
2974 )));
2975 }
2976
2977 let index_bytes = runmat_filesystem::read_async(&index_path)
2978 .await
2979 .map_err(|err| {
2980 data_error(format!(
2981 "failed to read chunk index '{}': {err}",
2982 index_path.display()
2983 ))
2984 })?;
2985 let mut chunk_index: DataChunkIndex = serde_json::from_slice(&index_bytes).map_err(|err| {
2986 data_error(format!(
2987 "failed to parse chunk index '{}': {err}",
2988 index_path.display()
2989 ))
2990 })?;
2991
2992 let mut pos_by_key = HashMap::new();
2993 for (idx, entry) in chunk_index.chunks.iter().enumerate() {
2994 pos_by_key.insert(entry.key.clone(), idx);
2995 }
2996 let touched = touched_chunk_coords(&ranges, &meta.chunk_shape, &meta.shape);
2997 let mut upload_batch = Vec::<(DataChunkDescriptor, Vec<u8>)>::new();
2998
2999 for coords in touched {
3000 let key = chunk_key(&coords);
3001 let chunk_start = chunk_start_for_coords(&coords, &meta.chunk_shape);
3002 let chunk_extent = chunk_extent_for_start(&chunk_start, &meta.chunk_shape, &meta.shape);
3003 let intersection = chunk_intersection(&ranges, &chunk_start, &chunk_extent);
3004 if intersection.is_empty() {
3005 continue;
3006 }
3007
3008 let (entry_index, existed, mut entry, mut chunk_payload) = load_or_init_chunk(
3009 root,
3010 array_name,
3011 &key,
3012 &coords,
3013 &chunk_extent,
3014 &meta.dtype,
3015 &pos_by_key,
3016 &chunk_index,
3017 )
3018 .await?;
3019 if chunk_payload.imaginary_values.is_none() && rhs.imaginary_values.is_some() {
3020 chunk_payload.imaginary_values = Some(crate::data::DataArrayValues::zeros(
3021 &meta.dtype,
3022 chunk_payload.values.len(),
3023 )?);
3024 } else if chunk_payload.imaginary_values.is_some() && rhs.imaginary_values.is_none() {
3025 return Err(data_error(
3026 "CLASS_MISMATCH: real and complex data-array slices cannot be assigned to one another",
3027 ));
3028 }
3029
3030 let mut local = vec![0usize; intersection.len()];
3031 let intersection_shape: Vec<usize> = intersection
3032 .iter()
3033 .map(|r| r.end.saturating_sub(r.start))
3034 .collect();
3035 loop {
3036 let mut global = Vec::with_capacity(intersection.len());
3037 for dim in 0..intersection.len() {
3038 global.push(intersection[dim].start + local[dim]);
3039 }
3040 let rhs_index: Vec<usize> = global
3041 .iter()
3042 .enumerate()
3043 .map(|(dim, g)| g.saturating_sub(ranges[dim].start))
3044 .collect();
3045 let chunk_index_local: Vec<usize> = global
3046 .iter()
3047 .enumerate()
3048 .map(|(dim, g)| g.saturating_sub(chunk_start[dim]))
3049 .collect();
3050 let rhs_linear = linear_index_column_major(&rhs_index, &rhs_shape)?;
3051 let chunk_linear = linear_index_column_major(&chunk_index_local, &chunk_extent)?;
3052 chunk_payload
3053 .values
3054 .set(chunk_linear, rhs.values.get(rhs_linear)?)?;
3055 if let (Some(target), Some(source)) =
3056 (&mut chunk_payload.imaginary_values, &rhs.imaginary_values)
3057 {
3058 target.set(chunk_linear, source.get(rhs_linear)?)?;
3059 }
3060 if !advance_index(&mut local, &intersection_shape) {
3061 break;
3062 }
3063 }
3064
3065 let chunk_bytes = serde_json::to_vec(&chunk_payload)
3066 .map_err(|err| data_error(format!("failed to encode chunk payload: {err}")))?;
3067 let chunk_path = root.join(&entry.data_path);
3068 runmat_filesystem::write_async(&chunk_path, &chunk_bytes)
3069 .await
3070 .map_err(|err| {
3071 data_error(format!(
3072 "failed to write chunk payload '{}': {err}",
3073 chunk_path.display()
3074 ))
3075 })?;
3076
3077 entry.coords = coords.clone();
3078 entry.shape = chunk_extent.clone();
3079 entry.bytes_raw = chunk_bytes.len() as u64;
3080 entry.bytes_stored = chunk_bytes.len() as u64;
3081 entry.hash = sha256_hex(&chunk_bytes);
3082 if existed {
3083 chunk_index.chunks[entry_index] = entry.clone();
3084 } else {
3085 chunk_index.chunks.push(entry.clone());
3086 pos_by_key.insert(key.clone(), chunk_index.chunks.len() - 1);
3087 }
3088 upload_batch.push((
3089 DataChunkDescriptor {
3090 key: key.clone(),
3091 object_id: entry.object_id.clone(),
3092 hash: entry.hash.clone(),
3093 bytes_raw: entry.bytes_raw,
3094 bytes_stored: entry.bytes_stored,
3095 },
3096 chunk_bytes,
3097 ));
3098 }
3099
3100 maybe_upload_chunk_batch_async(root, array_name, upload_batch).await?;
3101 tracing::info!(
3102 target: "runmat.data",
3103 dataset = %root.display(),
3104 array = array_name,
3105 touched_chunks = chunk_index.chunks.len(),
3106 "chunked slice write committed"
3107 );
3108 let index_write = serde_json::to_vec(&chunk_index)
3109 .map_err(|err| data_error(format!("failed to encode chunk index json: {err}")))?;
3110 runmat_filesystem::write_async(&index_path, &index_write)
3111 .await
3112 .map_err(|err| {
3113 data_error(format!(
3114 "failed to write chunk index '{}': {err}",
3115 index_path.display()
3116 ))
3117 })?;
3118
3119 if let Some(updated) = manifest.arrays.get_mut(array_name) {
3120 updated.shape = meta.shape.clone();
3121 updated.chunk_index_path = Some(index_rel_path.clone());
3122 }
3123 Ok(true)
3124}
3125
3126#[derive(Clone, Copy, Debug)]
3127struct DimRange {
3128 start: usize,
3129 end: usize,
3130}
3131
3132fn read_slice_payload(
3133 payload: &DataArrayPayload,
3134 slice_spec: &Value,
3135) -> BuiltinResult<DataArrayPayload> {
3136 let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
3137 let out_shape: Vec<usize> = ranges
3138 .iter()
3139 .map(|r| r.end.saturating_sub(r.start))
3140 .collect();
3141 let mut out_values = crate::data::DataArrayValues::zeros(&payload.dtype, 0)?;
3142 let mut out_imaginary_values = payload
3143 .imaginary_values
3144 .as_ref()
3145 .map(|_| crate::data::DataArrayValues::zeros(&payload.dtype, 0))
3146 .transpose()?;
3147 let mut out_index = vec![0usize; out_shape.len()];
3148 loop {
3149 let source_index: Vec<usize> = out_index
3150 .iter()
3151 .enumerate()
3152 .map(|(dim, idx)| ranges[dim].start + *idx)
3153 .collect();
3154 let linear = linear_index_column_major(&source_index, &payload.shape)?;
3155 out_values.push(payload.values.get(linear)?)?;
3156 if let (Some(source), Some(target)) = (&payload.imaginary_values, &mut out_imaginary_values)
3157 {
3158 target.push(source.get(linear)?)?;
3159 }
3160
3161 if !advance_index(&mut out_index, &out_shape) {
3162 break;
3163 }
3164 }
3165 Ok(DataArrayPayload {
3166 dtype: payload.dtype.clone(),
3167 shape: out_shape,
3168 values: out_values,
3169 imaginary_values: out_imaginary_values,
3170 })
3171}
3172
3173fn write_slice_payload(
3174 payload: &DataArrayPayload,
3175 slice_spec: &Value,
3176 rhs: &Value,
3177) -> BuiltinResult<DataArrayPayload> {
3178 let ranges = parse_slice_spec(slice_spec, &payload.shape)?;
3179 let target_shape: Vec<usize> = ranges
3180 .iter()
3181 .map(|r| r.end.saturating_sub(r.start))
3182 .collect();
3183 let rhs = DataArrayPayload::from_value(payload.dtype.clone(), rhs)?;
3184 if rhs.shape != target_shape {
3185 return Err(data_error(format!(
3186 "SHAPE_MISMATCH: rhs shape {:?} must match target slice shape {:?}",
3187 rhs.shape, target_shape
3188 )));
3189 }
3190
3191 let mut next = payload.values.clone();
3192 if payload.imaginary_values.is_some() != rhs.imaginary_values.is_some() {
3193 return Err(data_error(
3194 "CLASS_MISMATCH: real and complex data-array slices cannot be assigned to one another",
3195 ));
3196 }
3197 let mut next_imaginary = payload.imaginary_values.clone();
3198 let mut rhs_index = vec![0usize; target_shape.len()];
3199 let mut rhs_linear = 0usize;
3200 loop {
3201 let target_index: Vec<usize> = rhs_index
3202 .iter()
3203 .enumerate()
3204 .map(|(dim, idx)| ranges[dim].start + *idx)
3205 .collect();
3206 let target_linear = linear_index_column_major(&target_index, &payload.shape)?;
3207 next.set(target_linear, rhs.values.get(rhs_linear)?)?;
3208 if let (Some(target), Some(source)) = (&mut next_imaginary, &rhs.imaginary_values) {
3209 target.set(target_linear, source.get(rhs_linear)?)?;
3210 }
3211 rhs_linear += 1;
3212
3213 if !advance_index(&mut rhs_index, &target_shape) {
3214 break;
3215 }
3216 }
3217
3218 Ok(DataArrayPayload {
3219 dtype: payload.dtype.clone(),
3220 shape: payload.shape.clone(),
3221 values: next,
3222 imaginary_values: next_imaginary,
3223 })
3224}
3225
3226fn parse_slice_spec(slice_spec: &Value, shape: &[usize]) -> BuiltinResult<Vec<DimRange>> {
3227 match slice_spec {
3228 Value::Cell(cell) => {
3229 if cell.data.is_empty() {
3230 return Err(data_error("INVALID_SLICE: empty slice specification"));
3231 }
3232 let mut ranges = Vec::with_capacity(shape.len());
3233 for (dim, extent) in shape.iter().enumerate() {
3234 if let Some(item) = cell.data.get(dim) {
3235 ranges.push(parse_dim_range(item, *extent)?);
3236 } else {
3237 ranges.push(DimRange {
3238 start: 0,
3239 end: *extent,
3240 });
3241 }
3242 }
3243 Ok(ranges)
3244 }
3245 Value::String(s) if s == ":" => Ok(shape
3246 .iter()
3247 .map(|extent| DimRange {
3248 start: 0,
3249 end: *extent,
3250 })
3251 .collect()),
3252 _ => Err(data_error(
3253 "INVALID_SLICE: slice must be a cell spec like {1:10, :} or ':'",
3254 )),
3255 }
3256}
3257
3258fn parse_dim_range(value: &Value, extent: usize) -> BuiltinResult<DimRange> {
3259 if extent == 0 {
3260 return Ok(DimRange { start: 0, end: 0 });
3261 }
3262 match value {
3263 Value::String(s) if s == ":" => Ok(DimRange {
3264 start: 0,
3265 end: extent,
3266 }),
3267 Value::Num(n) => {
3268 let idx = floating_dimension_to_usize(*n)
3269 .and_then(|index| index.checked_sub(1))
3270 .ok_or_else(|| data_error("INVALID_SLICE: index out of bounds"))?;
3271 if idx >= extent {
3272 return Err(data_error("INVALID_SLICE: index out of bounds"));
3273 }
3274 Ok(DimRange {
3275 start: idx,
3276 end: idx + 1,
3277 })
3278 }
3279 Value::Int(i) => {
3280 let Some(idx) = i.try_to_usize().and_then(|index| index.checked_sub(1)) else {
3281 return Err(data_error("INVALID_SLICE: index out of bounds"));
3282 };
3283 if idx >= extent {
3284 return Err(data_error("INVALID_SLICE: index out of bounds"));
3285 }
3286 Ok(DimRange {
3287 start: idx,
3288 end: idx + 1,
3289 })
3290 }
3291 Value::Tensor(t) if tensor_utils::tensor_element_len(t) == 2 => {
3292 let start = t
3293 .numeric_value_at(0)
3294 .and_then(numeric_dimension_to_usize)
3295 .and_then(|index| index.checked_sub(1))
3296 .ok_or_else(|| data_error("INVALID_SLICE: range out of bounds"))?;
3297 let end_inclusive = t
3298 .numeric_value_at(1)
3299 .and_then(numeric_dimension_to_usize)
3300 .and_then(|index| index.checked_sub(1))
3301 .ok_or_else(|| data_error("INVALID_SLICE: range out of bounds"))?;
3302 if end_inclusive < start || end_inclusive >= extent {
3303 return Err(data_error("INVALID_SLICE: range out of bounds"));
3304 }
3305 Ok(DimRange {
3306 start,
3307 end: end_inclusive + 1,
3308 })
3309 }
3310 _ => Err(data_error(
3311 "INVALID_SLICE: dimension must be ':', scalar index, or [start end] range",
3312 )),
3313 }
3314}
3315
3316fn linear_index_column_major(index: &[usize], shape: &[usize]) -> BuiltinResult<usize> {
3317 if index.len() != shape.len() {
3318 return Err(data_error("INVALID_SLICE: rank mismatch"));
3319 }
3320 let mut stride = 1usize;
3321 let mut linear = 0usize;
3322 for (idx, extent) in index.iter().zip(shape.iter()) {
3323 if *idx >= *extent {
3324 return Err(data_error("INVALID_SLICE: index out of bounds"));
3325 }
3326 linear += idx * stride;
3327 stride = stride.saturating_mul(*extent);
3328 }
3329 Ok(linear)
3330}
3331
3332fn advance_index(index: &mut [usize], shape: &[usize]) -> bool {
3333 if shape.is_empty() {
3334 return false;
3335 }
3336 for dim in 0..shape.len() {
3337 index[dim] += 1;
3338 if index[dim] < shape[dim] {
3339 return true;
3340 }
3341 index[dim] = 0;
3342 }
3343 false
3344}
3345
3346fn chunk_key(coords: &[usize]) -> String {
3347 coords
3348 .iter()
3349 .map(|v| v.to_string())
3350 .collect::<Vec<_>>()
3351 .join(".")
3352}
3353
3354fn chunk_start_for_coords(coords: &[usize], chunk_shape: &[usize]) -> Vec<usize> {
3355 coords
3356 .iter()
3357 .enumerate()
3358 .map(|(dim, coord)| coord * chunk_shape.get(dim).copied().unwrap_or(1).max(1))
3359 .collect()
3360}
3361
3362fn chunk_extent_for_start(start: &[usize], chunk_shape: &[usize], shape: &[usize]) -> Vec<usize> {
3363 start
3364 .iter()
3365 .enumerate()
3366 .map(|(dim, start)| {
3367 let chunk = chunk_shape.get(dim).copied().unwrap_or(1).max(1);
3368 let end = (*start + chunk).min(shape[dim]);
3369 end.saturating_sub(*start)
3370 })
3371 .collect()
3372}
3373
3374fn chunk_intersection(
3375 ranges: &[DimRange],
3376 chunk_start: &[usize],
3377 chunk_extent: &[usize],
3378) -> Vec<DimRange> {
3379 let mut out = Vec::with_capacity(ranges.len());
3380 for dim in 0..ranges.len() {
3381 let c_start = chunk_start[dim];
3382 let c_end = c_start + chunk_extent[dim];
3383 let start = ranges[dim].start.max(c_start);
3384 let end = ranges[dim].end.min(c_end);
3385 if start >= end {
3386 return Vec::new();
3387 }
3388 out.push(DimRange { start, end });
3389 }
3390 out
3391}
3392
3393fn touched_chunk_coords(
3394 ranges: &[DimRange],
3395 chunk_shape: &[usize],
3396 shape: &[usize],
3397) -> Vec<Vec<usize>> {
3398 let mut span = Vec::with_capacity(ranges.len());
3399 let mut begin = Vec::with_capacity(ranges.len());
3400 for dim in 0..ranges.len() {
3401 if shape[dim] == 0 {
3402 return Vec::new();
3403 }
3404 let chunk = chunk_shape.get(dim).copied().unwrap_or(1).max(1);
3405 let first = ranges[dim].start / chunk;
3406 let last = (ranges[dim].end.saturating_sub(1)) / chunk;
3407 begin.push(first);
3408 span.push(last.saturating_sub(first) + 1);
3409 }
3410 let mut local = vec![0usize; span.len()];
3411 let mut out = Vec::new();
3412 loop {
3413 out.push(
3414 local
3415 .iter()
3416 .enumerate()
3417 .map(|(dim, v)| begin[dim] + *v)
3418 .collect::<Vec<_>>(),
3419 );
3420 if !advance_index(&mut local, &span) {
3421 break;
3422 }
3423 }
3424 out
3425}
3426
3427async fn maybe_upload_chunk_batch_async(
3428 root: &std::path::Path,
3429 array_name: &str,
3430 batch: Vec<(DataChunkDescriptor, Vec<u8>)>,
3431) -> BuiltinResult<()> {
3432 if batch.is_empty() {
3433 return Ok(());
3434 }
3435 let request = DataChunkUploadRequest {
3436 dataset_path: root.to_string_lossy().to_string(),
3437 array: array_name.to_string(),
3438 chunks: batch.iter().map(|(d, _)| d.clone()).collect(),
3439 };
3440 let targets = match runmat_filesystem::data_chunk_upload_targets_async(&request).await {
3441 Ok(targets) => targets,
3442 Err(err) if err.kind() == std::io::ErrorKind::Unsupported => return Ok(()),
3443 Err(err) => {
3444 return Err(data_error(format!(
3445 "failed to request data chunk upload targets: {err}"
3446 )))
3447 }
3448 };
3449 for (descriptor, bytes) in batch {
3450 let target = targets
3451 .iter()
3452 .find(|t| t.key == descriptor.key)
3453 .ok_or_else(|| {
3454 data_error(format!(
3455 "missing upload target for chunk '{}'",
3456 descriptor.key
3457 ))
3458 })?;
3459 runmat_filesystem::data_upload_chunk_async(target, &bytes)
3460 .await
3461 .map_err(|err| {
3462 data_error(format!(
3463 "failed to upload chunk '{}': {err}",
3464 descriptor.key
3465 ))
3466 })?;
3467 tracing::info!(
3468 target: "runmat.data",
3469 dataset = %root.display(),
3470 array = array_name,
3471 chunk_key = descriptor.key,
3472 bytes = bytes.len(),
3473 "chunk upload completed"
3474 );
3475 }
3476 Ok(())
3477}
3478
3479fn chunk_rel_path(array_name: &str, object_id: &str) -> String {
3480 format!("arrays/{array_name}/chunks/{object_id}.json")
3481}
3482
3483async fn load_or_init_chunk(
3484 root: &std::path::Path,
3485 array_name: &str,
3486 key: &str,
3487 coords: &[usize],
3488 chunk_extent: &[usize],
3489 dtype: &str,
3490 pos_by_key: &HashMap<String, usize>,
3491 chunk_index: &DataChunkIndex,
3492) -> BuiltinResult<(usize, bool, DataChunkIndexEntry, DataArrayPayload)> {
3493 if let Some(index) = pos_by_key.get(key).copied() {
3494 let entry = chunk_index
3495 .chunks
3496 .get(index)
3497 .cloned()
3498 .ok_or_else(|| data_error(format!("chunk index missing key '{key}'")))?;
3499 let bytes = runmat_filesystem::read_async(root.join(&entry.data_path))
3500 .await
3501 .map_err(|err| {
3502 data_error(format!(
3503 "failed to read chunk payload '{}': {err}",
3504 entry.data_path
3505 ))
3506 })?;
3507 let payload: DataArrayPayload = serde_json::from_slice::<DataArrayPayload>(&bytes)
3508 .map_err(|err| {
3509 data_error(format!(
3510 "failed to parse chunk payload '{}': {err}",
3511 entry.data_path
3512 ))
3513 })?
3514 .normalize_for_dtype(dtype)?;
3515 return Ok((index, true, entry, payload));
3516 }
3517
3518 let object_id = format!("obj_{}", key.replace('.', "_"));
3519 let entry = DataChunkIndexEntry {
3520 key: key.to_string(),
3521 object_id: object_id.clone(),
3522 hash: String::new(),
3523 bytes_raw: 0,
3524 bytes_stored: 0,
3525 coords: coords.to_vec(),
3526 shape: chunk_extent.to_vec(),
3527 data_path: chunk_rel_path(array_name, &object_id),
3528 };
3529 let payload = DataArrayPayload::zeros(dtype.to_string(), chunk_extent.to_vec())?;
3530 Ok((chunk_index.chunks.len(), false, entry, payload))
3531}
3532
3533fn attrs_to_struct(attrs: &BTreeMap<String, serde_json::Value>) -> Value {
3534 let mut out = StructValue::new();
3535 for (k, v) in attrs {
3536 out.fields.insert(k.clone(), json_to_value(v));
3537 }
3538 Value::Struct(out)
3539}
3540
3541fn value_to_json(value: &Value) -> serde_json::Value {
3542 match value {
3543 Value::String(s) => serde_json::Value::String(s.clone()),
3544 Value::CharArray(chars) => serde_json::Value::String(chars.data.iter().collect::<String>()),
3545 Value::Num(n) => serde_json::json!(n),
3546 Value::Int(i) => int_value_to_json(i),
3547 Value::Bool(b) => serde_json::json!(b),
3548 _ => serde_json::Value::String(format!("{value:?}")),
3549 }
3550}
3551
3552fn json_to_value(value: &serde_json::Value) -> Value {
3553 match value {
3554 serde_json::Value::Bool(b) => Value::Bool(*b),
3555 serde_json::Value::Number(n) => {
3556 if let Some(value) = n.as_i64() {
3557 Value::Int(IntValue::I64(value))
3558 } else if let Some(value) = n.as_u64() {
3559 Value::Int(IntValue::U64(value))
3560 } else {
3561 Value::Num(n.as_f64().unwrap_or_default())
3562 }
3563 }
3564 serde_json::Value::String(s) => Value::String(s.clone()),
3565 serde_json::Value::Array(arr) => {
3566 let vals = arr.iter().map(json_to_value).collect::<Vec<_>>();
3567 crate::make_cell(vals.clone(), 1, vals.len())
3568 .unwrap_or_else(|_| Value::String("<invalid-array>".to_string()))
3569 }
3570 serde_json::Value::Object(map) => {
3571 let mut s = StructValue::new();
3572 for (k, v) in map {
3573 s.fields.insert(k.clone(), json_to_value(v));
3574 }
3575 Value::Struct(s)
3576 }
3577 serde_json::Value::Null => Value::String("".to_string()),
3578 }
3579}
3580
3581#[cfg(all(test, not(target_arch = "wasm32")))]
3582mod tests {
3583 use super::*;
3584 use runmat_value::IntValue;
3585 use runmat_value::{ComplexTensor, IntegerComplexStorage};
3586
3587 fn tensor_values(tensor: &Tensor) -> Vec<f64> {
3588 tensor.materialize_f64()
3589 }
3590
3591 #[test]
3592 fn typed_data_slice_indices_do_not_saturate_through_i64() {
3593 let extent = usize::MAX;
3594 let result = parse_dim_range(&Value::Int(IntValue::U64(u64::MAX)), extent);
3595 match usize::try_from(u64::MAX)
3596 .ok()
3597 .and_then(|value| value.checked_sub(1))
3598 {
3599 Some(index) if index < extent => {
3600 let range = result.expect("index");
3601 assert_eq!(range.start, index);
3602 assert_eq!(range.end, index + 1);
3603 }
3604 _ => assert!(result.is_err()),
3605 }
3606 assert!(parse_dim_range(&Value::Int(IntValue::I64(-1)), extent).is_err());
3607 }
3608
3609 #[test]
3610 fn data_shape_and_slice_parsers_read_typed_integer_storage_exactly() {
3611 let cases = [
3612 runmat_value::IntegerStorage::I8(vec![2, 3]),
3613 runmat_value::IntegerStorage::I16(vec![2, 3]),
3614 runmat_value::IntegerStorage::I32(vec![2, 3]),
3615 runmat_value::IntegerStorage::I64(vec![2, 3]),
3616 runmat_value::IntegerStorage::U8(vec![2, 3]),
3617 runmat_value::IntegerStorage::U16(vec![2, 3]),
3618 runmat_value::IntegerStorage::U32(vec![2, 3]),
3619 runmat_value::IntegerStorage::U64(vec![2, 3]),
3620 ];
3621 for storage in cases {
3622 let shape = Tensor::new_integer(storage.clone(), vec![1, 2]).expect("shape");
3623 assert_eq!(
3624 parse_shape_from_value(&Value::Tensor(shape)).expect("shape"),
3625 vec![2, 3]
3626 );
3627
3628 let range = Tensor::new_integer(storage, vec![1, 2]).expect("range");
3629 let parsed = parse_dim_range(&Value::Tensor(range), 5).expect("range");
3630 assert_eq!(parsed.start, 1);
3631 assert_eq!(parsed.end, 3);
3632 }
3633
3634 let negative = Tensor::new_integer(runmat_value::IntegerStorage::I16(vec![-1]), vec![1, 1])
3635 .expect("shape");
3636 assert!(parse_shape_from_value(&Value::Tensor(negative)).is_err());
3637
3638 let invalid = Tensor::new_integer(
3639 runmat_value::IntegerStorage::U64(vec![u64::MAX, u64::MAX]),
3640 vec![1, 2],
3641 )
3642 .expect("range");
3643 assert!(parse_dim_range(&Value::Tensor(invalid), 5).is_err());
3644 }
3645
3646 #[test]
3647 fn data_shape_and_slice_floating_parsers_reject_fractional_and_unrepresentable_bounds() {
3648 assert!(parse_shape_from_value(&Value::Num(1.5)).is_err());
3649 let boundary_shape = parse_shape_from_value(&Value::Num(usize::MAX as f64));
3650 if usize::BITS == 64 {
3651 assert!(boundary_shape.is_err());
3652 } else {
3653 assert_eq!(boundary_shape.unwrap(), vec![usize::MAX]);
3654 }
3655 assert!(parse_shape_from_value(&Value::Num((usize::MAX as f64) + 1.0)).is_err());
3656
3657 let fractional_shape = Tensor::new(vec![2.0, 3.25], vec![1, 2]).expect("shape");
3658 assert!(parse_shape_from_value(&Value::Tensor(fractional_shape)).is_err());
3659 let single_shape = Tensor::from_f32(vec![2.0, 3.0], vec![1, 2]).expect("single shape");
3660 assert_eq!(
3661 parse_shape_from_value(&Value::Tensor(single_shape)).expect("single shape"),
3662 vec![2, 3]
3663 );
3664 let fractional_single_shape =
3665 Tensor::from_f32(vec![2.0, 3.25], vec![1, 2]).expect("fractional single shape");
3666 assert!(parse_shape_from_value(&Value::Tensor(fractional_single_shape)).is_err());
3667
3668 assert!(parse_dim_range(&Value::Num(1.5), 5).is_err());
3669 let boundary_range = parse_dim_range(&Value::Num(usize::MAX as f64), usize::MAX);
3670 if usize::BITS == 64 {
3671 assert!(boundary_range.is_err());
3672 } else {
3673 let parsed = boundary_range.expect("32-bit boundary");
3674 assert_eq!(parsed.start, usize::MAX - 1);
3675 assert_eq!(parsed.end, usize::MAX);
3676 }
3677
3678 let fractional_range = Tensor::new(vec![1.0, 3.5], vec![1, 2]).expect("range");
3679 assert!(parse_dim_range(&Value::Tensor(fractional_range), 5).is_err());
3680 let single_range = Tensor::from_f32(vec![2.0, 3.0], vec![1, 2]).expect("single range");
3681 let parsed = parse_dim_range(&Value::Tensor(single_range), 5).expect("single range");
3682 assert_eq!((parsed.start, parsed.end), (1, 3));
3683
3684 let too_large_range = Tensor::new(vec![1.0, usize::MAX as f64], vec![1, 2]).expect("range");
3685 let tensor_boundary = parse_dim_range(&Value::Tensor(too_large_range), usize::MAX);
3686 if usize::BITS == 64 {
3687 assert!(tensor_boundary.is_err());
3688 } else {
3689 let parsed = tensor_boundary.expect("32-bit tensor boundary");
3690 assert_eq!(parsed.start, 0);
3691 assert_eq!(parsed.end, usize::MAX);
3692 }
3693 }
3694 use crate::dispatcher::call_builtin;
3695 use async_trait::async_trait;
3696 use axum::extract::{Query, State};
3697 use axum::http::{HeaderMap, StatusCode};
3698 use axum::routing::{post, put};
3699 use axum::{Json, Router};
3700 use runmat_filesystem::data_contract::{
3701 DataChunkUploadRequest, DataChunkUploadTarget, DataManifestDescriptor, DataManifestRequest,
3702 };
3703 use runmat_filesystem::{
3704 DirEntry, FileHandle, FsMetadata, FsProvider, NativeFsProvider, OpenFlags,
3705 };
3706 use runmat_value::{CellArray, IntegerStorage};
3707 use serde::Deserialize;
3708 use std::path::Path;
3709 use std::sync::{Arc, Mutex, MutexGuard};
3710 use tokio::runtime::Runtime;
3711 use tokio::sync::oneshot;
3712
3713 fn serial_test_guard() -> MutexGuard<'static, ()> {
3714 runmat_filesystem::provider_override_lock()
3715 }
3716
3717 fn native_provider_guard() -> runmat_filesystem::ProviderGuard {
3718 runmat_filesystem::replace_provider(Arc::new(NativeFsProvider))
3719 }
3720
3721 #[test]
3722 fn attribute_json_preserves_native_integer_text() {
3723 for value in [
3724 runmat_value::IntValue::I64(i64::MIN),
3725 runmat_value::IntValue::U64(u64::MAX),
3726 ] {
3727 let json = value_to_json(&Value::Int(value.clone()));
3728 assert_eq!(json.to_string(), value.decimal_string());
3729 }
3730 }
3731
3732 #[test]
3733 fn attribute_json_roundtrips_wide_integer_values_without_f64() {
3734 for value in [IntValue::I64(i64::MIN), IntValue::U64(u64::MAX)] {
3735 let json = value_to_json(&Value::Int(value.clone()));
3736 assert_eq!(json_to_value(&json), Value::Int(value));
3737 }
3738 }
3739
3740 #[test]
3741 fn io_data_descriptors_cover_constructor_and_transaction_surface() {
3742 let data_labels: Vec<&str> = DATA_CREATE_DESCRIPTOR
3743 .signatures
3744 .iter()
3745 .map(|sig| sig.label)
3746 .collect();
3747 assert!(data_labels.contains(&"ds = data.create(path, schema)"));
3748
3749 let read_labels: Vec<&str> = DATAARRAY_READ_DESCRIPTOR
3750 .signatures
3751 .iter()
3752 .map(|sig| sig.label)
3753 .collect();
3754 assert!(read_labels.contains(&"X = DataArray.read(arr, sliceSpec)"));
3755
3756 let write_labels: Vec<&str> = DATAARRAY_WRITE_DESCRIPTOR
3757 .signatures
3758 .iter()
3759 .map(|sig| sig.label)
3760 .collect();
3761 assert!(write_labels.contains(&"tf = DataArray.write(arr, values)"));
3762 assert!(write_labels.contains(&"tf = DataArray.write(arr, sliceSpec, values)"));
3763
3764 let tx_labels: Vec<&str> = DATATX_COMMIT_DESCRIPTOR
3765 .signatures
3766 .iter()
3767 .map(|sig| sig.label)
3768 .collect();
3769 assert!(tx_labels.contains(&"tf = DataTransaction.commit(tx, options)"));
3770
3771 assert_eq!(DATA_CREATE_INTEGER_CAPABILITIES.len(), 1);
3772 assert_eq!(DATAARRAY_READ_INTEGER_CAPABILITIES.len(), 1);
3773 assert_eq!(DATAARRAY_WRITE_INTEGER_CAPABILITIES.len(), 1);
3774 assert_eq!(DATAARRAY_RESIZE_INTEGER_CAPABILITIES.len(), 1);
3775 assert_eq!(DATAARRAY_FILL_INTEGER_CAPABILITIES.len(), 1);
3776 assert_eq!(DATASET_ATTRS_INTEGER_CAPABILITIES.len(), 1);
3777 assert_eq!(DATASET_GET_ATTR_INTEGER_CAPABILITIES.len(), 2);
3778 assert_eq!(DATASET_SET_ATTR_INTEGER_CAPABILITIES.len(), 1);
3779 assert_eq!(DATASET_SET_ATTRS_INTEGER_CAPABILITIES.len(), 1);
3780 assert_eq!(DATATX_WRITE_INTEGER_CAPABILITIES.len(), 1);
3781 assert_eq!(DATATX_SET_ATTR_INTEGER_CAPABILITIES.len(), 1);
3782 assert_eq!(DATATX_SET_ATTRS_INTEGER_CAPABILITIES.len(), 1);
3783 assert_eq!(DATATX_RESIZE_INTEGER_CAPABILITIES.len(), 1);
3784 assert_eq!(DATATX_FILL_INTEGER_CAPABILITIES.len(), 1);
3785 assert_eq!(DATATX_CREATE_ARRAY_INTEGER_CAPABILITIES.len(), 1);
3786 for capability in [
3787 &DATA_CREATE_INTEGER_CAPABILITIES[0],
3788 &DATAARRAY_READ_INTEGER_CAPABILITIES[0],
3789 &DATAARRAY_WRITE_INTEGER_CAPABILITIES[0],
3790 &DATAARRAY_RESIZE_INTEGER_CAPABILITIES[0],
3791 &DATAARRAY_FILL_INTEGER_CAPABILITIES[0],
3792 &DATASET_ATTRS_INTEGER_CAPABILITIES[0],
3793 &DATASET_GET_ATTR_INTEGER_CAPABILITIES[0],
3794 &DATASET_GET_ATTR_INTEGER_CAPABILITIES[1],
3795 &DATASET_SET_ATTR_INTEGER_CAPABILITIES[0],
3796 &DATASET_SET_ATTRS_INTEGER_CAPABILITIES[0],
3797 &DATATX_WRITE_INTEGER_CAPABILITIES[0],
3798 &DATATX_SET_ATTR_INTEGER_CAPABILITIES[0],
3799 &DATATX_SET_ATTRS_INTEGER_CAPABILITIES[0],
3800 &DATATX_RESIZE_INTEGER_CAPABILITIES[0],
3801 &DATATX_FILL_INTEGER_CAPABILITIES[0],
3802 &DATATX_CREATE_ARRAY_INTEGER_CAPABILITIES[0],
3803 ] {
3804 assert!(capability.inputs.iter().all(|input| input.classes
3805 == crate::builtins::common::integer_capability::ALL_INTEGER_CLASSES));
3806 }
3807 }
3808
3809 #[test]
3810 fn data_create_uses_exact_all_class_schema_controls_and_integer_payloads() {
3811 let _serial = serial_test_guard();
3812 let _provider = native_provider_guard();
3813 let cases = vec![
3814 (
3815 "int8",
3816 IntegerStorage::I8(vec![2, 3]),
3817 IntegerStorage::I8(vec![1, 3]),
3818 IntegerStorage::I8(vec![0; 6]),
3819 ),
3820 (
3821 "int16",
3822 IntegerStorage::I16(vec![2, 3]),
3823 IntegerStorage::I16(vec![1, 3]),
3824 IntegerStorage::I16(vec![0; 6]),
3825 ),
3826 (
3827 "int32",
3828 IntegerStorage::I32(vec![2, 3]),
3829 IntegerStorage::I32(vec![1, 3]),
3830 IntegerStorage::I32(vec![0; 6]),
3831 ),
3832 (
3833 "int64",
3834 IntegerStorage::I64(vec![2, 3]),
3835 IntegerStorage::I64(vec![1, 3]),
3836 IntegerStorage::I64(vec![0; 6]),
3837 ),
3838 (
3839 "uint8",
3840 IntegerStorage::U8(vec![2, 3]),
3841 IntegerStorage::U8(vec![1, 3]),
3842 IntegerStorage::U8(vec![0; 6]),
3843 ),
3844 (
3845 "uint16",
3846 IntegerStorage::U16(vec![2, 3]),
3847 IntegerStorage::U16(vec![1, 3]),
3848 IntegerStorage::U16(vec![0; 6]),
3849 ),
3850 (
3851 "uint32",
3852 IntegerStorage::U32(vec![2, 3]),
3853 IntegerStorage::U32(vec![1, 3]),
3854 IntegerStorage::U32(vec![0; 6]),
3855 ),
3856 (
3857 "uint64",
3858 IntegerStorage::U64(vec![2, 3]),
3859 IntegerStorage::U64(vec![1, 3]),
3860 IntegerStorage::U64(vec![0; 6]),
3861 ),
3862 ];
3863
3864 for (dtype, shape, chunk, expected) in cases {
3865 let dir = tempfile::tempdir().expect("tempdir");
3866 let path = dir
3867 .path()
3868 .join(format!("schema-{dtype}.data"))
3869 .to_string_lossy()
3870 .to_string();
3871 let mut meta = StructValue::new();
3872 meta.fields
3873 .insert("dtype".to_string(), Value::String(dtype.to_string()));
3874 meta.fields.insert(
3875 "shape".to_string(),
3876 Value::Tensor(Tensor::new_integer(shape, vec![1, 2]).expect("shape")),
3877 );
3878 meta.fields.insert(
3879 "chunk".to_string(),
3880 Value::Tensor(Tensor::new_integer(chunk, vec![1, 2]).expect("chunk")),
3881 );
3882 let mut arrays = StructValue::new();
3883 arrays
3884 .fields
3885 .insert("samples".to_string(), Value::Struct(meta));
3886 let mut schema = StructValue::new();
3887 schema
3888 .fields
3889 .insert("arrays".to_string(), Value::Struct(arrays));
3890
3891 let dataset =
3892 call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
3893 .unwrap_or_else(|error| panic!("{dtype}: create failed: {error}"));
3894 let array = call_builtin(
3895 "Dataset.array",
3896 &[dataset, Value::String("samples".to_string())],
3897 )
3898 .expect("array handle");
3899 let Value::Tensor(values) =
3900 call_builtin("DataArray.read", &[array]).expect("read zeros")
3901 else {
3902 panic!("{dtype}: expected tensor");
3903 };
3904 assert_eq!(values.shape, vec![2, 3], "{dtype}");
3905 assert_eq!(values.integer_storage(), Some(&expected), "{dtype}");
3906 }
3907 }
3908
3909 #[test]
3910 fn data_namespace_lifecycle_rejects_integer_arguments_and_unsupported_extras() {
3911 let integer_values = [
3912 IntValue::I8(1),
3913 IntValue::I16(1),
3914 IntValue::I32(1),
3915 IntValue::I64(1),
3916 IntValue::U8(1),
3917 IntValue::U16(1),
3918 IntValue::U32(1),
3919 IntValue::U64(1),
3920 ];
3921 for integer in integer_values {
3922 let value = Value::Int(integer);
3923 for (name, args) in [
3924 (
3925 "data.create",
3926 vec![value.clone(), Value::Struct(StructValue::new())],
3927 ),
3928 ("data.open", vec![value.clone()]),
3929 ("data.exists", vec![value.clone()]),
3930 ("data.delete", vec![value.clone()]),
3931 ("data.inspect", vec![value.clone()]),
3932 ("data.list", vec![value.clone()]),
3933 (
3934 "data.copy",
3935 vec![value.clone(), Value::String("to.data".to_string())],
3936 ),
3937 (
3938 "data.copy",
3939 vec![Value::String("from.data".to_string()), value.clone()],
3940 ),
3941 (
3942 "data.move",
3943 vec![value.clone(), Value::String("to.data".to_string())],
3944 ),
3945 (
3946 "data.move",
3947 vec![Value::String("from.data".to_string()), value.clone()],
3948 ),
3949 (
3950 "data.import",
3951 vec![
3952 value.clone(),
3953 Value::String("data".to_string()),
3954 Value::String("source.data".to_string()),
3955 ],
3956 ),
3957 (
3958 "data.import",
3959 vec![
3960 Value::String("path.data".to_string()),
3961 value.clone(),
3962 Value::String("source.data".to_string()),
3963 ],
3964 ),
3965 (
3966 "data.import",
3967 vec![
3968 Value::String("path.data".to_string()),
3969 Value::String("data".to_string()),
3970 value.clone(),
3971 ],
3972 ),
3973 (
3974 "data.export",
3975 vec![
3976 value.clone(),
3977 Value::String("data".to_string()),
3978 Value::String("target.data".to_string()),
3979 ],
3980 ),
3981 (
3982 "data.export",
3983 vec![
3984 Value::String("path.data".to_string()),
3985 value.clone(),
3986 Value::String("target.data".to_string()),
3987 ],
3988 ),
3989 (
3990 "data.export",
3991 vec![
3992 Value::String("path.data".to_string()),
3993 Value::String("data".to_string()),
3994 value.clone(),
3995 ],
3996 ),
3997 ("commit", vec![value]),
3998 ] {
3999 let error = call_builtin(name, &args).expect_err("expected integer rejection");
4000 assert!(!error.message().is_empty(), "{name}");
4001 }
4002 }
4003
4004 let error = call_builtin(
4005 "data.exists",
4006 &[
4007 Value::String("missing.data".to_string()),
4008 Value::Int(IntValue::U8(1)),
4009 ],
4010 )
4011 .expect_err("fixed documented arity must reject extras");
4012 assert!(!error.message().is_empty());
4013
4014 let error = call_builtin(
4015 "DataTransaction.commit",
4016 &[Value::Int(IntValue::U8(1)), Value::Int(IntValue::U8(1))],
4017 )
4018 .expect_err("commit options must be a struct");
4019 assert!(error.message().contains("options must be a scalar struct"));
4020
4021 let error = call_builtin(
4022 "DataTransaction.commit",
4023 &[
4024 Value::Int(IntValue::U8(1)),
4025 Value::Struct(StructValue::new()),
4026 Value::Struct(StructValue::new()),
4027 ],
4028 )
4029 .expect_err("commit must reject more than one options struct");
4030 assert!(error.message().contains("at most one options struct"));
4031 }
4032
4033 #[test]
4034 fn dataset_nonattribute_methods_reject_integer_receivers() {
4035 for (name, extra) in [
4036 ("Dataset.array", Some(Value::String("samples".to_string()))),
4037 ("Dataset.arrays", None),
4038 ("Dataset.begin", None),
4039 (
4040 "Dataset.has_array",
4041 Some(Value::String("samples".to_string())),
4042 ),
4043 ("Dataset.id", None),
4044 ("Dataset.path", None),
4045 ("Dataset.refresh", None),
4046 (
4047 "Dataset.snapshot",
4048 Some(Value::String("checkpoint".to_string())),
4049 ),
4050 ("Dataset.version", None),
4051 ] {
4052 let mut args = vec![Value::Int(IntValue::U64(u64::MAX))];
4053 args.extend(extra);
4054 let error = call_builtin(name, &args).expect_err("method requires a Dataset receiver");
4055 assert!(
4056 error.message().contains("expected object"),
4057 "{name}: {error}"
4058 );
4059 }
4060 }
4061
4062 #[test]
4063 fn dataset_attributes_keep_exact_integer_values_and_documented_class_boundaries() {
4064 let _serial = serial_test_guard();
4065 let _provider = native_provider_guard();
4066 let dir = tempfile::tempdir().expect("tempdir");
4067 let path = dir
4068 .path()
4069 .join("integer-attrs.data")
4070 .to_string_lossy()
4071 .to_string();
4072 let mut array_meta = StructValue::new();
4073 array_meta
4074 .fields
4075 .insert("dtype".to_string(), Value::String("f64".to_string()));
4076 array_meta.fields.insert(
4077 "shape".to_string(),
4078 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("shape")),
4079 );
4080 let mut arrays = StructValue::new();
4081 arrays
4082 .fields
4083 .insert("base".to_string(), Value::Struct(array_meta));
4084 let mut schema = StructValue::new();
4085 schema
4086 .fields
4087 .insert("arrays".to_string(), Value::Struct(arrays));
4088 let ds = call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
4089 .expect("create dataset");
4090 let cases = [
4091 ("int8", IntValue::I8(i8::MIN), IntValue::I64(i8::MIN as i64)),
4092 (
4093 "int16",
4094 IntValue::I16(i16::MIN),
4095 IntValue::I64(i16::MIN as i64),
4096 ),
4097 (
4098 "int32",
4099 IntValue::I32(i32::MIN),
4100 IntValue::I64(i32::MIN as i64),
4101 ),
4102 ("int64", IntValue::I64(i64::MIN), IntValue::I64(i64::MIN)),
4103 (
4104 "uint8",
4105 IntValue::U8(u8::MAX),
4106 IntValue::I64(u8::MAX as i64),
4107 ),
4108 (
4109 "uint16",
4110 IntValue::U16(u16::MAX),
4111 IntValue::I64(u16::MAX as i64),
4112 ),
4113 (
4114 "uint32",
4115 IntValue::U32(u32::MAX),
4116 IntValue::I64(u32::MAX as i64),
4117 ),
4118 ("uint64", IntValue::U64(u64::MAX), IntValue::U64(u64::MAX)),
4119 ];
4120
4121 let mut bulk = StructValue::new();
4122 for (name, source, _) in &cases {
4123 call_builtin(
4124 "Dataset.set_attr",
4125 &[
4126 ds.clone(),
4127 Value::String(format!("direct_{name}")),
4128 Value::Int(source.clone()),
4129 ],
4130 )
4131 .expect("set integer attribute");
4132 bulk.fields
4133 .insert(format!("bulk_{name}"), Value::Int(source.clone()));
4134 }
4135 call_builtin("Dataset.set_attrs", &[ds.clone(), Value::Struct(bulk)])
4136 .expect("set integer attributes");
4137
4138 let Value::Struct(attrs) =
4139 call_builtin("Dataset.attrs", std::slice::from_ref(&ds)).expect("read attributes")
4140 else {
4141 panic!("expected attribute struct");
4142 };
4143 for (name, source, canonical) in cases {
4144 let expected = Value::Int(canonical);
4145 assert_eq!(
4146 call_builtin(
4147 "Dataset.get_attr",
4148 &[ds.clone(), Value::String(format!("direct_{name}"))],
4149 )
4150 .expect("get integer attribute"),
4151 expected,
4152 "{name} direct"
4153 );
4154 assert_eq!(
4155 attrs.fields.get(&format!("bulk_{name}")),
4156 Some(&expected),
4157 "{name} bulk"
4158 );
4159 assert_eq!(
4160 call_builtin(
4161 "Dataset.get_attr",
4162 &[
4163 ds.clone(),
4164 Value::String(format!("missing_{name}")),
4165 Value::Int(source.clone()),
4166 ],
4167 )
4168 .expect("get integer default"),
4169 Value::Int(source),
4170 "{name} default"
4171 );
4172 }
4173 }
4174
4175 #[test]
4176 fn data_transaction_lifecycle_methods_reject_integer_receivers() {
4177 for (name, extra) in [
4178 ("DataTransaction.abort", None),
4179 ("DataTransaction.commit", None),
4180 (
4181 "DataTransaction.delete_array",
4182 Some(Value::String("samples".to_string())),
4183 ),
4184 ("DataTransaction.id", None),
4185 ("DataTransaction.status", None),
4186 ] {
4187 let mut args = vec![Value::Int(IntValue::U64(u64::MAX))];
4188 args.extend(extra);
4189 let error = call_builtin(name, &args)
4190 .expect_err("lifecycle method requires a DataTransaction receiver");
4191 assert!(
4192 error.message().contains("expected object"),
4193 "{name}: {error}"
4194 );
4195 }
4196 }
4197
4198 #[test]
4199 fn data_array_metadata_accessors_reject_integer_receivers() {
4200 for name in [
4201 "DataArray.chunk_shape",
4202 "DataArray.codec",
4203 "DataArray.dtype",
4204 "DataArray.name",
4205 "DataArray.rank",
4206 "DataArray.shape",
4207 ] {
4208 let error = call_builtin(name, &[Value::Int(IntValue::U64(u64::MAX))])
4209 .expect_err("metadata accessor requires a DataArray receiver");
4210 assert!(
4211 error.message().contains("expected object"),
4212 "{name}: {error}"
4213 );
4214 }
4215 }
4216
4217 #[derive(Default)]
4218 struct CountingDataUploadProvider {
4219 inner: NativeFsProvider,
4220 uploaded_keys: Arc<Mutex<Vec<String>>>,
4221 }
4222
4223 struct HttpDataUploadProvider {
4224 inner: NativeFsProvider,
4225 base_url: String,
4226 client: reqwest::blocking::Client,
4227 }
4228
4229 impl HttpDataUploadProvider {
4230 fn new(base_url: String) -> Self {
4231 Self {
4232 inner: NativeFsProvider,
4233 base_url,
4234 client: reqwest::blocking::Client::new(),
4235 }
4236 }
4237 }
4238
4239 #[async_trait(?Send)]
4240 impl FsProvider for HttpDataUploadProvider {
4241 fn open(&self, path: &Path, flags: &OpenFlags) -> std::io::Result<Box<dyn FileHandle>> {
4242 self.inner.open(path, flags)
4243 }
4244
4245 async fn read(&self, path: &Path) -> std::io::Result<Vec<u8>> {
4246 self.inner.read(path).await
4247 }
4248
4249 async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
4250 self.inner.write(path, data).await
4251 }
4252
4253 async fn remove_file(&self, path: &Path) -> std::io::Result<()> {
4254 self.inner.remove_file(path).await
4255 }
4256
4257 async fn metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
4258 self.inner.metadata(path).await
4259 }
4260
4261 async fn symlink_metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
4262 self.inner.symlink_metadata(path).await
4263 }
4264
4265 async fn read_dir(&self, path: &Path) -> std::io::Result<Vec<DirEntry>> {
4266 self.inner.read_dir(path).await
4267 }
4268
4269 async fn canonicalize(&self, path: &Path) -> std::io::Result<std::path::PathBuf> {
4270 self.inner.canonicalize(path).await
4271 }
4272
4273 async fn create_dir(&self, path: &Path) -> std::io::Result<()> {
4274 self.inner.create_dir(path).await
4275 }
4276
4277 async fn create_dir_all(&self, path: &Path) -> std::io::Result<()> {
4278 self.inner.create_dir_all(path).await
4279 }
4280
4281 async fn remove_dir(&self, path: &Path) -> std::io::Result<()> {
4282 self.inner.remove_dir(path).await
4283 }
4284
4285 async fn remove_dir_all(&self, path: &Path) -> std::io::Result<()> {
4286 self.inner.remove_dir_all(path).await
4287 }
4288
4289 async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
4290 self.inner.rename(from, to).await
4291 }
4292
4293 async fn set_readonly(&self, path: &Path, readonly: bool) -> std::io::Result<()> {
4294 self.inner.set_readonly(path, readonly).await
4295 }
4296
4297 async fn data_manifest_descriptor(
4298 &self,
4299 request: &DataManifestRequest,
4300 ) -> std::io::Result<DataManifestDescriptor> {
4301 self.inner.data_manifest_descriptor(request).await
4302 }
4303
4304 async fn data_chunk_upload_targets(
4305 &self,
4306 request: &DataChunkUploadRequest,
4307 ) -> std::io::Result<Vec<DataChunkUploadTarget>> {
4308 #[derive(Deserialize)]
4309 struct UploadTargetsResponse {
4310 targets: Vec<DataChunkUploadTarget>,
4311 }
4312 let url = format!("{}/data/chunks/upload-targets", self.base_url);
4313 let response = self
4314 .client
4315 .post(url)
4316 .json(request)
4317 .send()
4318 .map_err(|err| std::io::Error::other(err.to_string()))?;
4319 if !response.status().is_success() {
4320 return Err(std::io::Error::other(format!(
4321 "upload targets request failed: {}",
4322 response.status()
4323 )));
4324 }
4325 let parsed: UploadTargetsResponse = response
4326 .json()
4327 .map_err(|err| std::io::Error::other(err.to_string()))?;
4328 Ok(parsed.targets)
4329 }
4330
4331 async fn data_upload_chunk(
4332 &self,
4333 target: &DataChunkUploadTarget,
4334 data: &[u8],
4335 ) -> std::io::Result<()> {
4336 let upload_url = if let Some(key) = target.upload_url.strip_prefix("upload://") {
4337 format!("{}/upload?key={}", self.base_url, key)
4338 } else {
4339 target.upload_url.clone()
4340 };
4341 let method = reqwest::Method::from_bytes(target.method.as_bytes())
4342 .map_err(|err| std::io::Error::other(err.to_string()))?;
4343 let mut request = self.client.request(method, &upload_url);
4344 for (k, v) in &target.headers {
4345 request = request.header(k, v);
4346 }
4347 let response = request
4348 .body(data.to_vec())
4349 .send()
4350 .map_err(|err| std::io::Error::other(err.to_string()))?;
4351 if !response.status().is_success() {
4352 return Err(std::io::Error::other(format!(
4353 "chunk upload failed: {}",
4354 response.status()
4355 )));
4356 }
4357 Ok(())
4358 }
4359 }
4360
4361 #[derive(Clone, Default)]
4362 struct UploadHarness {
4363 uploads: Arc<Mutex<Vec<String>>>,
4364 }
4365
4366 #[derive(Deserialize)]
4367 struct UploadChunkQuery {
4368 key: String,
4369 }
4370
4371 async fn upload_targets_handler(
4372 Json(req): Json<DataChunkUploadRequest>,
4373 ) -> Result<Json<serde_json::Value>, StatusCode> {
4374 let targets = req
4375 .chunks
4376 .iter()
4377 .map(|chunk| {
4378 serde_json::json!({
4379 "key": chunk.key,
4380 "method": "PUT",
4381 "upload_url": format!("upload://{}", chunk.key),
4382 "headers": {
4383 "x-runmat-hash": chunk.hash,
4384 }
4385 })
4386 })
4387 .collect::<Vec<_>>();
4388 Ok(Json(serde_json::json!({ "targets": targets })))
4389 }
4390
4391 async fn upload_handler(
4392 State(harness): State<UploadHarness>,
4393 Query(query): Query<UploadChunkQuery>,
4394 headers: HeaderMap,
4395 body: axum::body::Bytes,
4396 ) -> Result<(), StatusCode> {
4397 if body.is_empty() {
4398 return Err(StatusCode::BAD_REQUEST);
4399 }
4400 if headers.get("x-runmat-hash").is_none() {
4401 return Err(StatusCode::BAD_REQUEST);
4402 }
4403 let mut guard = harness.uploads.lock().expect("uploads lock poisoned");
4404 guard.push(query.key);
4405 Ok(())
4406 }
4407
4408 fn spawn_upload_server() -> (
4409 String,
4410 Arc<Mutex<Vec<String>>>,
4411 Runtime,
4412 oneshot::Sender<()>,
4413 ) {
4414 let harness = UploadHarness::default();
4415 let uploads = Arc::clone(&harness.uploads);
4416 let runtime = Runtime::new().expect("tokio runtime");
4417 let (addr, shutdown_tx) = runtime.block_on(async move {
4418 let listener = tokio::net::TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
4419 .await
4420 .expect("bind upload server");
4421 let addr = listener.local_addr().expect("local addr");
4422 let app = Router::new()
4423 .route("/data/chunks/upload-targets", post(upload_targets_handler))
4424 .route("/upload", put(upload_handler))
4425 .with_state(harness);
4426 let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>();
4427 let server = axum::serve(listener, app).with_graceful_shutdown(async {
4428 let _ = shutdown_rx.await;
4429 });
4430 tokio::spawn(async move {
4431 let _ = server.await;
4432 });
4433 (addr, shutdown_tx)
4434 });
4435 (format!("http://{}", addr), uploads, runtime, shutdown_tx)
4436 }
4437
4438 impl CountingDataUploadProvider {
4439 fn uploaded_keys(&self) -> Arc<Mutex<Vec<String>>> {
4440 Arc::clone(&self.uploaded_keys)
4441 }
4442 }
4443
4444 #[async_trait(?Send)]
4445 impl FsProvider for CountingDataUploadProvider {
4446 fn open(&self, path: &Path, flags: &OpenFlags) -> std::io::Result<Box<dyn FileHandle>> {
4447 self.inner.open(path, flags)
4448 }
4449
4450 async fn read(&self, path: &Path) -> std::io::Result<Vec<u8>> {
4451 self.inner.read(path).await
4452 }
4453
4454 async fn write(&self, path: &Path, data: &[u8]) -> std::io::Result<()> {
4455 self.inner.write(path, data).await
4456 }
4457
4458 async fn remove_file(&self, path: &Path) -> std::io::Result<()> {
4459 self.inner.remove_file(path).await
4460 }
4461
4462 async fn metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
4463 self.inner.metadata(path).await
4464 }
4465
4466 async fn symlink_metadata(&self, path: &Path) -> std::io::Result<FsMetadata> {
4467 self.inner.symlink_metadata(path).await
4468 }
4469
4470 async fn read_dir(&self, path: &Path) -> std::io::Result<Vec<DirEntry>> {
4471 self.inner.read_dir(path).await
4472 }
4473
4474 async fn canonicalize(&self, path: &Path) -> std::io::Result<std::path::PathBuf> {
4475 self.inner.canonicalize(path).await
4476 }
4477
4478 async fn create_dir(&self, path: &Path) -> std::io::Result<()> {
4479 self.inner.create_dir(path).await
4480 }
4481
4482 async fn create_dir_all(&self, path: &Path) -> std::io::Result<()> {
4483 self.inner.create_dir_all(path).await
4484 }
4485
4486 async fn remove_dir(&self, path: &Path) -> std::io::Result<()> {
4487 self.inner.remove_dir(path).await
4488 }
4489
4490 async fn remove_dir_all(&self, path: &Path) -> std::io::Result<()> {
4491 self.inner.remove_dir_all(path).await
4492 }
4493
4494 async fn rename(&self, from: &Path, to: &Path) -> std::io::Result<()> {
4495 self.inner.rename(from, to).await
4496 }
4497
4498 async fn set_readonly(&self, path: &Path, readonly: bool) -> std::io::Result<()> {
4499 self.inner.set_readonly(path, readonly).await
4500 }
4501
4502 async fn data_manifest_descriptor(
4503 &self,
4504 request: &DataManifestRequest,
4505 ) -> std::io::Result<DataManifestDescriptor> {
4506 self.inner.data_manifest_descriptor(request).await
4507 }
4508
4509 async fn data_chunk_upload_targets(
4510 &self,
4511 request: &DataChunkUploadRequest,
4512 ) -> std::io::Result<Vec<DataChunkUploadTarget>> {
4513 Ok(request
4514 .chunks
4515 .iter()
4516 .map(|chunk| DataChunkUploadTarget {
4517 key: chunk.key.clone(),
4518 method: "PUT".to_string(),
4519 upload_url: format!("count://{}", chunk.object_id),
4520 headers: std::collections::HashMap::new(),
4521 })
4522 .collect())
4523 }
4524
4525 async fn data_upload_chunk(
4526 &self,
4527 target: &DataChunkUploadTarget,
4528 _data: &[u8],
4529 ) -> std::io::Result<()> {
4530 let mut guard = match self.uploaded_keys.lock() {
4531 Ok(guard) => guard,
4532 Err(poisoned) => poisoned.into_inner(),
4533 };
4534 guard.push(target.key.clone());
4535 Ok(())
4536 }
4537 }
4538
4539 #[test]
4540 fn create_open_write_read_dataset() {
4541 let _serial = serial_test_guard();
4542 let _provider_guard = native_provider_guard();
4543 let dir = tempfile::tempdir().expect("tempdir");
4544 let path = dir.path().join("sample.data").to_string_lossy().to_string();
4545
4546 let mut array_meta = StructValue::new();
4547 array_meta
4548 .fields
4549 .insert("dtype".to_string(), Value::String("f64".to_string()));
4550 array_meta.fields.insert(
4551 "shape".to_string(),
4552 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape tensor")),
4553 );
4554 let mut arrays = StructValue::new();
4555 arrays
4556 .fields
4557 .insert("temperature".to_string(), Value::Struct(array_meta));
4558 let mut schema = StructValue::new();
4559 schema
4560 .fields
4561 .insert("arrays".to_string(), Value::Struct(arrays));
4562
4563 let ds = call_builtin(
4564 "data.create",
4565 &[Value::String(path.clone()), Value::Struct(schema)],
4566 )
4567 .expect("create dataset");
4568
4569 let arr = call_builtin(
4570 "Dataset.array",
4571 &[ds, Value::String("temperature".to_string())],
4572 )
4573 .expect("dataset array");
4574 let write_tensor = Tensor::new(vec![1.0, 2.0, 3.0, 4.0], vec![2, 2]).expect("write tensor");
4575 call_builtin(
4576 "DataArray.write",
4577 &[arr.clone(), Value::Tensor(write_tensor)],
4578 )
4579 .expect("write array");
4580
4581 let read_back = call_builtin("DataArray.read", &[arr]).expect("read array");
4582 let Value::Tensor(t) = read_back else {
4583 panic!("expected tensor");
4584 };
4585 assert_eq!(t.shape, vec![2, 2]);
4586 assert_eq!(tensor_values(&t), vec![1.0, 2.0, 3.0, 4.0]);
4587 }
4588
4589 #[test]
4590 fn write_and_read_slice_payload() {
4591 let _serial = serial_test_guard();
4592 let _provider_guard = native_provider_guard();
4593 let dir = tempfile::tempdir().expect("tempdir");
4594 let path = dir.path().join("slice.data").to_string_lossy().to_string();
4595
4596 let mut array_meta = StructValue::new();
4597 array_meta
4598 .fields
4599 .insert("dtype".to_string(), Value::String("f64".to_string()));
4600 array_meta.fields.insert(
4601 "shape".to_string(),
4602 Value::Tensor(Tensor::new(vec![3.0, 3.0], vec![1, 2]).expect("shape tensor")),
4603 );
4604 let mut arrays = StructValue::new();
4605 arrays
4606 .fields
4607 .insert("temperature".to_string(), Value::Struct(array_meta));
4608 let mut schema = StructValue::new();
4609 schema
4610 .fields
4611 .insert("arrays".to_string(), Value::Struct(arrays));
4612
4613 let ds = call_builtin(
4614 "data.create",
4615 &[Value::String(path.clone()), Value::Struct(schema)],
4616 )
4617 .expect("create dataset");
4618 let arr = call_builtin(
4619 "Dataset.array",
4620 &[ds, Value::String("temperature".to_string())],
4621 )
4622 .expect("dataset array");
4623
4624 let slice = Value::Cell(
4625 CellArray::new(
4626 vec![
4627 Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
4628 Value::String(":".to_string()),
4629 ],
4630 1,
4631 2,
4632 )
4633 .expect("slice cell"),
4634 );
4635 let rhs = Value::Tensor(
4636 Tensor::new(vec![10.0, 11.0, 12.0, 13.0, 14.0, 15.0], vec![2, 3]).expect("rhs"),
4637 );
4638 call_builtin("DataArray.write", &[arr.clone(), slice.clone(), rhs]).expect("slice write");
4639
4640 let read_back = call_builtin("DataArray.read", &[arr.clone(), slice]).expect("slice read");
4641 let Value::Tensor(t) = read_back else {
4642 panic!("expected tensor");
4643 };
4644 assert_eq!(t.shape, vec![2, 3]);
4645 assert_eq!(tensor_values(&t), vec![10.0, 11.0, 12.0, 13.0, 14.0, 15.0]);
4646 }
4647
4648 #[test]
4649 fn slice_write_updates_only_touched_chunks() {
4650 let _serial = serial_test_guard();
4651 let _provider_guard = native_provider_guard();
4652 let dir = tempfile::tempdir().expect("tempdir");
4653 let path = dir
4654 .path()
4655 .join("chunked.data")
4656 .to_string_lossy()
4657 .to_string();
4658
4659 let mut array_meta = StructValue::new();
4660 array_meta
4661 .fields
4662 .insert("dtype".to_string(), Value::String("f64".to_string()));
4663 array_meta.fields.insert(
4664 "shape".to_string(),
4665 Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
4666 );
4667 array_meta.fields.insert(
4668 "chunk".to_string(),
4669 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
4670 );
4671 let mut arrays = StructValue::new();
4672 arrays
4673 .fields
4674 .insert("temperature".to_string(), Value::Struct(array_meta));
4675 let mut schema = StructValue::new();
4676 schema
4677 .fields
4678 .insert("arrays".to_string(), Value::Struct(arrays));
4679
4680 let ds = call_builtin(
4681 "data.create",
4682 &[Value::String(path.clone()), Value::Struct(schema)],
4683 )
4684 .expect("create dataset");
4685 let arr = call_builtin(
4686 "Dataset.array",
4687 &[ds, Value::String("temperature".to_string())],
4688 )
4689 .expect("dataset array");
4690
4691 let full = Value::Tensor(
4692 Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4]).expect("full tensor"),
4693 );
4694 call_builtin("DataArray.write", &[arr.clone(), full]).expect("initial write");
4695
4696 let root = std::path::PathBuf::from(&path);
4697 let untouched_path = root.join("arrays/temperature/chunks/obj_1_1.json");
4698 let touched_path = root.join("arrays/temperature/chunks/obj_0_0.json");
4699 let untouched_before =
4700 futures::executor::block_on(runmat_filesystem::read_async(&untouched_path))
4701 .expect("read untouched before");
4702 let touched_before =
4703 futures::executor::block_on(runmat_filesystem::read_async(&touched_path))
4704 .expect("read touched before");
4705
4706 let slice = Value::Cell(
4707 CellArray::new(
4708 vec![
4709 Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
4710 Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
4711 ],
4712 1,
4713 2,
4714 )
4715 .expect("slice cell"),
4716 );
4717 let rhs =
4718 Value::Tensor(Tensor::new(vec![99.0, 98.0, 97.0, 96.0], vec![2, 2]).expect("rhs"));
4719 call_builtin("DataArray.write", &[arr.clone(), slice, rhs]).expect("slice write");
4720
4721 let untouched_after =
4722 futures::executor::block_on(runmat_filesystem::read_async(&untouched_path))
4723 .expect("read untouched after");
4724 let touched_after =
4725 futures::executor::block_on(runmat_filesystem::read_async(&touched_path))
4726 .expect("read touched after");
4727 assert_eq!(untouched_before, untouched_after);
4728 assert_ne!(touched_before, touched_after);
4729 }
4730
4731 #[test]
4732 fn slice_write_uploads_only_touched_chunk_targets() {
4733 let _serial = serial_test_guard();
4734 let provider = Arc::new(CountingDataUploadProvider::default());
4735 let uploaded = provider.uploaded_keys();
4736 let _guard = runmat_filesystem::replace_provider(provider);
4737
4738 let dir = tempfile::tempdir().expect("tempdir");
4739 let path = dir
4740 .path()
4741 .join("remote-chunked.data")
4742 .to_string_lossy()
4743 .to_string();
4744
4745 let mut array_meta = StructValue::new();
4746 array_meta
4747 .fields
4748 .insert("dtype".to_string(), Value::String("f64".to_string()));
4749 array_meta.fields.insert(
4750 "shape".to_string(),
4751 Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
4752 );
4753 array_meta.fields.insert(
4754 "chunk".to_string(),
4755 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
4756 );
4757 let mut arrays = StructValue::new();
4758 arrays
4759 .fields
4760 .insert("temperature".to_string(), Value::Struct(array_meta));
4761 let mut schema = StructValue::new();
4762 schema
4763 .fields
4764 .insert("arrays".to_string(), Value::Struct(arrays));
4765
4766 let ds = call_builtin(
4767 "data.create",
4768 &[Value::String(path.clone()), Value::Struct(schema)],
4769 )
4770 .expect("create dataset");
4771 let arr = call_builtin(
4772 "Dataset.array",
4773 &[ds, Value::String("temperature".to_string())],
4774 )
4775 .expect("dataset array");
4776
4777 call_builtin(
4778 "DataArray.write",
4779 &[
4780 arr.clone(),
4781 Value::Tensor(
4782 Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
4783 .expect("full tensor"),
4784 ),
4785 ],
4786 )
4787 .expect("initial write");
4788
4789 let manifest =
4790 futures::executor::block_on(crate::data::read_manifest_async(&dataset_root(&path)))
4791 .expect("manifest after initial write");
4792 let meta = manifest
4793 .arrays
4794 .get("temperature")
4795 .expect("temperature meta");
4796 let chunk_index_path =
4797 dataset_root(&path).join(meta.chunk_index_path.clone().expect("chunk index path"));
4798 assert!(
4799 futures::executor::block_on(runmat_filesystem::metadata_async(&chunk_index_path))
4800 .is_ok()
4801 );
4802
4803 {
4804 let mut keys = uploaded.lock().expect("uploaded keys lock");
4805 keys.clear();
4806 }
4807
4808 let slice = Value::Cell(
4809 CellArray::new(
4810 vec![
4811 Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
4812 Value::Tensor(Tensor::new(vec![1.0, 2.0], vec![1, 2]).expect("range")),
4813 ],
4814 1,
4815 2,
4816 )
4817 .expect("slice cell"),
4818 );
4819 let rhs = Value::Tensor(Tensor::new(vec![9.0, 8.0, 7.0, 6.0], vec![2, 2]).expect("rhs"));
4820 call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
4821
4822 let keys = uploaded.lock().expect("uploaded keys lock");
4823 assert_eq!(keys.as_slice(), ["0.0".to_string()].as_slice());
4824 }
4825
4826 #[test]
4827 fn slice_write_uploads_expected_cross_boundary_chunk_targets() {
4828 let _serial = serial_test_guard();
4829 let provider = Arc::new(CountingDataUploadProvider::default());
4830 let uploaded = provider.uploaded_keys();
4831 let _guard = runmat_filesystem::replace_provider(provider);
4832
4833 let dir = tempfile::tempdir().expect("tempdir");
4834 let path = dir
4835 .path()
4836 .join("remote-chunked-boundary.data")
4837 .to_string_lossy()
4838 .to_string();
4839
4840 let mut array_meta = StructValue::new();
4841 array_meta
4842 .fields
4843 .insert("dtype".to_string(), Value::String("f64".to_string()));
4844 array_meta.fields.insert(
4845 "shape".to_string(),
4846 Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
4847 );
4848 array_meta.fields.insert(
4849 "chunk".to_string(),
4850 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
4851 );
4852 let mut arrays = StructValue::new();
4853 arrays
4854 .fields
4855 .insert("temperature".to_string(), Value::Struct(array_meta));
4856 let mut schema = StructValue::new();
4857 schema
4858 .fields
4859 .insert("arrays".to_string(), Value::Struct(arrays));
4860
4861 let ds = call_builtin(
4862 "data.create",
4863 &[Value::String(path.clone()), Value::Struct(schema)],
4864 )
4865 .expect("create dataset");
4866 let arr = call_builtin(
4867 "Dataset.array",
4868 &[ds, Value::String("temperature".to_string())],
4869 )
4870 .expect("dataset array");
4871
4872 call_builtin(
4873 "DataArray.write",
4874 &[
4875 arr.clone(),
4876 Value::Tensor(
4877 Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
4878 .expect("full tensor"),
4879 ),
4880 ],
4881 )
4882 .expect("initial write");
4883 {
4884 let mut keys = uploaded.lock().expect("uploaded keys lock");
4885 keys.clear();
4886 }
4887
4888 let slice = Value::Cell(
4889 CellArray::new(
4890 vec![
4891 Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
4892 Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
4893 ],
4894 1,
4895 2,
4896 )
4897 .expect("slice cell"),
4898 );
4899 let rhs =
4900 Value::Tensor(Tensor::new(vec![19.0, 18.0, 17.0, 16.0], vec![2, 2]).expect("rhs"));
4901 call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
4902
4903 let mut keys = uploaded.lock().expect("uploaded keys lock").clone();
4904 keys.sort();
4905 keys.dedup();
4906 assert_eq!(
4907 keys.as_slice(),
4908 [
4909 "0.0".to_string(),
4910 "0.1".to_string(),
4911 "1.0".to_string(),
4912 "1.1".to_string(),
4913 ]
4914 .as_slice()
4915 );
4916 }
4917
4918 #[test]
4919 fn slice_write_hits_http_server_data_endpoints_with_expected_keys() {
4920 let _serial = serial_test_guard();
4921 let (base_url, uploads, runtime, shutdown_tx) = spawn_upload_server();
4922 let provider = Arc::new(HttpDataUploadProvider::new(base_url));
4923 let _guard = runmat_filesystem::replace_provider(provider);
4924
4925 let dir = tempfile::tempdir().expect("tempdir");
4926 let path = dir
4927 .path()
4928 .join("http-endpoint.data")
4929 .to_string_lossy()
4930 .to_string();
4931
4932 let mut array_meta = StructValue::new();
4933 array_meta
4934 .fields
4935 .insert("dtype".to_string(), Value::String("f64".to_string()));
4936 array_meta.fields.insert(
4937 "shape".to_string(),
4938 Value::Tensor(Tensor::new(vec![4.0, 4.0], vec![1, 2]).expect("shape tensor")),
4939 );
4940 array_meta.fields.insert(
4941 "chunk".to_string(),
4942 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("chunk tensor")),
4943 );
4944 let mut arrays = StructValue::new();
4945 arrays
4946 .fields
4947 .insert("temperature".to_string(), Value::Struct(array_meta));
4948 let mut schema = StructValue::new();
4949 schema
4950 .fields
4951 .insert("arrays".to_string(), Value::Struct(arrays));
4952
4953 let ds = call_builtin(
4954 "data.create",
4955 &[Value::String(path.clone()), Value::Struct(schema)],
4956 )
4957 .expect("create dataset");
4958 let arr = call_builtin(
4959 "Dataset.array",
4960 &[ds, Value::String("temperature".to_string())],
4961 )
4962 .expect("dataset array");
4963
4964 call_builtin(
4965 "DataArray.write",
4966 &[
4967 arr.clone(),
4968 Value::Tensor(
4969 Tensor::new((1..=16).map(|v| v as f64).collect(), vec![4, 4])
4970 .expect("full tensor"),
4971 ),
4972 ],
4973 )
4974 .expect("initial write");
4975
4976 {
4977 let mut keys = uploads.lock().expect("uploads lock");
4978 keys.clear();
4979 }
4980
4981 let slice = Value::Cell(
4982 CellArray::new(
4983 vec![
4984 Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
4985 Value::Tensor(Tensor::new(vec![2.0, 3.0], vec![1, 2]).expect("range")),
4986 ],
4987 1,
4988 2,
4989 )
4990 .expect("slice cell"),
4991 );
4992 let rhs =
4993 Value::Tensor(Tensor::new(vec![19.0, 18.0, 17.0, 16.0], vec![2, 2]).expect("rhs"));
4994 call_builtin("DataArray.write", &[arr, slice, rhs]).expect("slice write");
4995
4996 let mut keys = uploads.lock().expect("uploads lock").clone();
4997 keys.sort();
4998 keys.dedup();
4999 assert_eq!(
5000 keys.as_slice(),
5001 [
5002 "0.0".to_string(),
5003 "0.1".to_string(),
5004 "1.0".to_string(),
5005 "1.1".to_string(),
5006 ]
5007 .as_slice()
5008 );
5009
5010 let _ = shutdown_tx.send(());
5011 drop(runtime);
5012 }
5013
5014 #[test]
5015 fn tx_create_resize_fill_and_delete_array() {
5016 let _serial = serial_test_guard();
5017 let _provider = native_provider_guard();
5018 let dir = tempfile::tempdir().expect("tempdir");
5019 let path = dir.path().join("tx-ops.data").to_string_lossy().to_string();
5020
5021 let mut arrays = StructValue::new();
5022 let mut array_meta = StructValue::new();
5023 array_meta
5024 .fields
5025 .insert("dtype".to_string(), Value::String("f64".to_string()));
5026 array_meta.fields.insert(
5027 "shape".to_string(),
5028 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("shape tensor")),
5029 );
5030 arrays
5031 .fields
5032 .insert("base".to_string(), Value::Struct(array_meta));
5033 let mut schema = StructValue::new();
5034 schema
5035 .fields
5036 .insert("arrays".to_string(), Value::Struct(arrays));
5037
5038 let ds = call_builtin(
5039 "data.create",
5040 &[Value::String(path.clone()), Value::Struct(schema)],
5041 )
5042 .expect("create dataset");
5043
5044 let tx = call_builtin("Dataset.begin", &[ds]).expect("begin tx");
5045 let mut new_meta = StructValue::new();
5046 new_meta
5047 .fields
5048 .insert("dtype".to_string(), Value::String("f64".to_string()));
5049 new_meta.fields.insert(
5050 "shape".to_string(),
5051 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape tensor")),
5052 );
5053 call_builtin(
5054 "DataTransaction.create_array",
5055 &[
5056 tx.clone(),
5057 Value::String("new_array".to_string()),
5058 Value::Struct(new_meta),
5059 ],
5060 )
5061 .expect("create array in tx");
5062 call_builtin(
5063 "DataTransaction.resize",
5064 &[
5065 tx.clone(),
5066 Value::String("new_array".to_string()),
5067 Value::Tensor(Tensor::new(vec![3.0, 1.0], vec![1, 2]).expect("shape tensor")),
5068 ],
5069 )
5070 .expect("resize array in tx");
5071 call_builtin(
5072 "DataTransaction.fill",
5073 &[
5074 tx.clone(),
5075 Value::String("new_array".to_string()),
5076 Value::Num(7.0),
5077 ],
5078 )
5079 .expect("fill array in tx");
5080 call_builtin(
5081 "DataTransaction.delete_array",
5082 &[tx.clone(), Value::String("base".to_string())],
5083 )
5084 .expect("delete array in tx");
5085 call_builtin("commit", &[tx, Value::Struct(StructValue::new())])
5086 .expect("commit tx through alias options form");
5087
5088 let ds = call_builtin("data.open", &[Value::String(path)]).expect("open dataset");
5089 let has_base = call_builtin(
5090 "Dataset.has_array",
5091 &[ds.clone(), Value::String("base".to_string())],
5092 )
5093 .expect("has base");
5094 assert_eq!(has_base, Value::Bool(false));
5095 let arr = call_builtin(
5096 "Dataset.array",
5097 &[ds, Value::String("new_array".to_string())],
5098 )
5099 .expect("new array");
5100 let read_back = call_builtin("DataArray.read", &[arr]).expect("read array");
5101 let Value::Tensor(t) = read_back else {
5102 panic!("expected tensor");
5103 };
5104 assert_eq!(t.shape, vec![3, 1]);
5105 assert_eq!(tensor_values(&t), vec![7.0, 7.0, 7.0]);
5106 }
5107
5108 #[test]
5109 fn data_transactions_preserve_every_integer_class_and_typed_controls() {
5110 let _serial = serial_test_guard();
5111 let _provider = native_provider_guard();
5112 let cases = vec![
5113 (
5114 "int8",
5115 IntegerStorage::I8(vec![i8::MIN, i8::MAX]),
5116 IntegerStorage::I8(vec![2, 1]),
5117 IntegerStorage::I8(vec![1, 1]),
5118 IntegerStorage::I8(vec![1, 2]),
5119 IntValue::I8(1),
5120 IntValue::I8(i8::MIN),
5121 ),
5122 (
5123 "int16",
5124 IntegerStorage::I16(vec![i16::MIN, i16::MAX]),
5125 IntegerStorage::I16(vec![2, 1]),
5126 IntegerStorage::I16(vec![1, 1]),
5127 IntegerStorage::I16(vec![1, 2]),
5128 IntValue::I16(1),
5129 IntValue::I16(i16::MIN),
5130 ),
5131 (
5132 "int32",
5133 IntegerStorage::I32(vec![i32::MIN, i32::MAX]),
5134 IntegerStorage::I32(vec![2, 1]),
5135 IntegerStorage::I32(vec![1, 1]),
5136 IntegerStorage::I32(vec![1, 2]),
5137 IntValue::I32(1),
5138 IntValue::I32(i32::MIN),
5139 ),
5140 (
5141 "int64",
5142 IntegerStorage::I64(vec![i64::MIN, i64::MAX]),
5143 IntegerStorage::I64(vec![2, 1]),
5144 IntegerStorage::I64(vec![1, 1]),
5145 IntegerStorage::I64(vec![1, 2]),
5146 IntValue::I64(1),
5147 IntValue::I64(i64::MIN),
5148 ),
5149 (
5150 "uint8",
5151 IntegerStorage::U8(vec![0, u8::MAX]),
5152 IntegerStorage::U8(vec![2, 1]),
5153 IntegerStorage::U8(vec![1, 1]),
5154 IntegerStorage::U8(vec![1, 2]),
5155 IntValue::U8(1),
5156 IntValue::U8(u8::MAX),
5157 ),
5158 (
5159 "uint16",
5160 IntegerStorage::U16(vec![0, u16::MAX]),
5161 IntegerStorage::U16(vec![2, 1]),
5162 IntegerStorage::U16(vec![1, 1]),
5163 IntegerStorage::U16(vec![1, 2]),
5164 IntValue::U16(1),
5165 IntValue::U16(u16::MAX),
5166 ),
5167 (
5168 "uint32",
5169 IntegerStorage::U32(vec![0, u32::MAX]),
5170 IntegerStorage::U32(vec![2, 1]),
5171 IntegerStorage::U32(vec![1, 1]),
5172 IntegerStorage::U32(vec![1, 2]),
5173 IntValue::U32(1),
5174 IntValue::U32(u32::MAX),
5175 ),
5176 (
5177 "uint64",
5178 IntegerStorage::U64(vec![0, u64::MAX]),
5179 IntegerStorage::U64(vec![2, 1]),
5180 IntegerStorage::U64(vec![1, 1]),
5181 IntegerStorage::U64(vec![1, 2]),
5182 IntValue::U64(1),
5183 IntValue::U64(u64::MAX),
5184 ),
5185 ];
5186
5187 for (dtype, values, shape, chunk, range, one, attribute) in cases {
5188 let dir = tempfile::tempdir().expect("tempdir");
5189 let path = dir
5190 .path()
5191 .join(format!("tx-{dtype}.data"))
5192 .to_string_lossy()
5193 .to_string();
5194 let mut base_meta = StructValue::new();
5195 base_meta
5196 .fields
5197 .insert("dtype".to_string(), Value::String("f64".to_string()));
5198 base_meta.fields.insert(
5199 "shape".to_string(),
5200 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("base shape")),
5201 );
5202 let mut arrays = StructValue::new();
5203 arrays
5204 .fields
5205 .insert("base".to_string(), Value::Struct(base_meta));
5206 let mut schema = StructValue::new();
5207 schema
5208 .fields
5209 .insert("arrays".to_string(), Value::Struct(arrays));
5210 let ds = call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
5211 .expect("create dataset");
5212 let tx = call_builtin("Dataset.begin", std::slice::from_ref(&ds))
5213 .expect("begin transaction");
5214
5215 let mut meta = StructValue::new();
5216 meta.fields
5217 .insert("dtype".to_string(), Value::String(dtype.to_string()));
5218 meta.fields.insert(
5219 "shape".to_string(),
5220 Value::Tensor(Tensor::new_integer(shape.clone(), vec![1, 2]).expect("shape")),
5221 );
5222 meta.fields.insert(
5223 "chunk".to_string(),
5224 Value::Tensor(Tensor::new_integer(chunk, vec![1, 2]).expect("chunk")),
5225 );
5226 call_builtin(
5227 "DataTransaction.create_array",
5228 &[
5229 tx.clone(),
5230 Value::String("typed".to_string()),
5231 Value::Struct(meta),
5232 ],
5233 )
5234 .expect("queue create");
5235 call_builtin(
5236 "DataTransaction.resize",
5237 &[
5238 tx.clone(),
5239 Value::String("typed".to_string()),
5240 Value::Tensor(Tensor::new_integer(shape, vec![1, 2]).expect("resize shape")),
5241 ],
5242 )
5243 .expect("queue resize");
5244 call_builtin(
5245 "DataTransaction.fill",
5246 &[
5247 tx.clone(),
5248 Value::String("typed".to_string()),
5249 Value::Int(attribute.clone()),
5250 ],
5251 )
5252 .expect("queue fill");
5253 let slice = Value::Cell(
5254 CellArray::new(
5255 vec![
5256 Value::Tensor(Tensor::new_integer(range, vec![1, 2]).expect("row range")),
5257 Value::Int(one),
5258 ],
5259 1,
5260 2,
5261 )
5262 .expect("slice"),
5263 );
5264 call_builtin(
5265 "DataTransaction.write",
5266 &[
5267 tx.clone(),
5268 Value::String("typed".to_string()),
5269 slice,
5270 Value::Tensor(Tensor::new_integer(values.clone(), vec![2, 1]).expect("values")),
5271 ],
5272 )
5273 .expect("queue write");
5274 call_builtin(
5275 "DataTransaction.set_attr",
5276 &[
5277 tx.clone(),
5278 Value::String(format!("{dtype}_one")),
5279 Value::Int(attribute.clone()),
5280 ],
5281 )
5282 .expect("queue attribute");
5283 let mut attrs = StructValue::new();
5284 attrs
5285 .fields
5286 .insert(format!("{dtype}_many"), Value::Int(attribute));
5287 call_builtin(
5288 "DataTransaction.set_attrs",
5289 &[tx.clone(), Value::Struct(attrs)],
5290 )
5291 .expect("queue attributes");
5292 call_builtin("DataTransaction.commit", &[tx]).expect("commit transaction");
5293
5294 let arr = call_builtin("Dataset.array", &[ds, Value::String("typed".to_string())])
5295 .expect("open typed array");
5296 let Value::Tensor(read_back) =
5297 call_builtin("DataArray.read", &[arr]).expect("read typed array")
5298 else {
5299 panic!("expected tensor");
5300 };
5301 assert_eq!(read_back.integer_storage(), Some(&values), "{dtype}");
5302 }
5303 }
5304
5305 #[test]
5306 fn data_arrays_preserve_every_integer_class_through_chunked_write_and_read() {
5307 let _serial = serial_test_guard();
5308 let _provider = native_provider_guard();
5309 let cases = vec![
5310 (
5311 "int8",
5312 IntegerStorage::I8(vec![i8::MIN, i8::MAX]),
5313 IntValue::I8(i8::MIN),
5314 IntegerStorage::I8(vec![i8::MIN; 2]),
5315 ),
5316 (
5317 "int16",
5318 IntegerStorage::I16(vec![i16::MIN, i16::MAX]),
5319 IntValue::I16(i16::MIN),
5320 IntegerStorage::I16(vec![i16::MIN; 2]),
5321 ),
5322 (
5323 "int32",
5324 IntegerStorage::I32(vec![i32::MIN, i32::MAX]),
5325 IntValue::I32(i32::MIN),
5326 IntegerStorage::I32(vec![i32::MIN; 2]),
5327 ),
5328 (
5329 "int64",
5330 IntegerStorage::I64(vec![i64::MIN, i64::MAX]),
5331 IntValue::I64(i64::MIN),
5332 IntegerStorage::I64(vec![i64::MIN; 2]),
5333 ),
5334 (
5335 "uint8",
5336 IntegerStorage::U8(vec![0, u8::MAX]),
5337 IntValue::U8(u8::MAX),
5338 IntegerStorage::U8(vec![u8::MAX; 2]),
5339 ),
5340 (
5341 "uint16",
5342 IntegerStorage::U16(vec![0, u16::MAX]),
5343 IntValue::U16(u16::MAX),
5344 IntegerStorage::U16(vec![u16::MAX; 2]),
5345 ),
5346 (
5347 "uint32",
5348 IntegerStorage::U32(vec![0, u32::MAX]),
5349 IntValue::U32(u32::MAX),
5350 IntegerStorage::U32(vec![u32::MAX; 2]),
5351 ),
5352 (
5353 "uint64",
5354 IntegerStorage::U64(vec![0, u64::MAX]),
5355 IntValue::U64(u64::MAX),
5356 IntegerStorage::U64(vec![u64::MAX; 2]),
5357 ),
5358 ];
5359
5360 for (dtype, storage, fill, filled_storage) in cases {
5361 let dir = tempfile::tempdir().expect("tempdir");
5362 let path = dir
5363 .path()
5364 .join(format!("{dtype}.data"))
5365 .to_string_lossy()
5366 .to_string();
5367 let mut array_meta = StructValue::new();
5368 array_meta
5369 .fields
5370 .insert("dtype".to_string(), Value::String(dtype.to_string()));
5371 array_meta.fields.insert(
5372 "shape".to_string(),
5373 Value::Tensor(Tensor::new(vec![2.0, 1.0], vec![1, 2]).expect("shape")),
5374 );
5375 array_meta.fields.insert(
5376 "chunk".to_string(),
5377 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("chunk")),
5378 );
5379 let mut arrays = StructValue::new();
5380 arrays
5381 .fields
5382 .insert("samples".to_string(), Value::Struct(array_meta));
5383 let mut schema = StructValue::new();
5384 schema
5385 .fields
5386 .insert("arrays".to_string(), Value::Struct(arrays));
5387
5388 let ds = call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
5389 .expect("create dataset");
5390 let arr = call_builtin("Dataset.array", &[ds, Value::String("samples".to_string())])
5391 .expect("array");
5392 let input = Tensor::new_integer(storage.clone(), vec![2, 1]).expect("integer tensor");
5393 call_builtin("DataArray.write", &[arr.clone(), Value::Tensor(input)])
5394 .expect("write integer array");
5395
5396 let Value::Tensor(read_back) =
5397 call_builtin("DataArray.read", std::slice::from_ref(&arr)).expect("read")
5398 else {
5399 panic!("expected tensor");
5400 };
5401 assert_eq!(read_back.integer_storage(), Some(&storage), "{dtype}");
5402
5403 call_builtin("DataArray.fill", &[arr.clone(), Value::Int(fill)])
5404 .expect("fill integer array");
5405 let Value::Tensor(read_back) =
5406 call_builtin("DataArray.read", &[arr]).expect("read filled array")
5407 else {
5408 panic!("expected tensor");
5409 };
5410 assert_eq!(
5411 read_back.integer_storage(),
5412 Some(&filled_storage),
5413 "{dtype} fill"
5414 );
5415 }
5416 }
5417
5418 #[test]
5419 fn uint64_data_array_slice_fill_and_transaction_paths_remain_exact() {
5420 let _serial = serial_test_guard();
5421 let _provider = native_provider_guard();
5422 let dir = tempfile::tempdir().expect("tempdir");
5423 let path = dir.path().join("uint64.data").to_string_lossy().to_string();
5424 let mut array_meta = StructValue::new();
5425 array_meta
5426 .fields
5427 .insert("dtype".to_string(), Value::String("uint64".to_string()));
5428 array_meta.fields.insert(
5429 "shape".to_string(),
5430 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape")),
5431 );
5432 array_meta.fields.insert(
5433 "chunk".to_string(),
5434 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("chunk")),
5435 );
5436 let mut arrays = StructValue::new();
5437 arrays
5438 .fields
5439 .insert("samples".to_string(), Value::Struct(array_meta));
5440 let mut schema = StructValue::new();
5441 schema
5442 .fields
5443 .insert("arrays".to_string(), Value::Struct(arrays));
5444 let ds = call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
5445 .expect("create dataset");
5446 let arr = call_builtin(
5447 "Dataset.array",
5448 &[ds.clone(), Value::String("samples".to_string())],
5449 )
5450 .expect("array");
5451
5452 call_builtin(
5453 "DataArray.fill",
5454 &[
5455 arr.clone(),
5456 Value::Int(runmat_value::IntValue::U64(u64::MAX)),
5457 ],
5458 )
5459 .expect("fill");
5460 let slice = Value::Cell(
5461 CellArray::new(
5462 vec![
5463 Value::Int(runmat_value::IntValue::I32(1)),
5464 Value::String(":".to_string()),
5465 ],
5466 1,
5467 2,
5468 )
5469 .expect("slice"),
5470 );
5471 let replacement = Tensor::new_integer(
5472 IntegerStorage::U64(vec![1_u64 << 63, u64::MAX - 1]),
5473 vec![1, 2],
5474 )
5475 .expect("replacement");
5476 call_builtin(
5477 "DataArray.write",
5478 &[arr.clone(), slice, Value::Tensor(replacement)],
5479 )
5480 .expect("slice write");
5481
5482 let tx = call_builtin("Dataset.begin", &[ds]).expect("begin transaction");
5483 call_builtin(
5484 "DataTransaction.fill",
5485 &[
5486 tx.clone(),
5487 Value::String("samples".to_string()),
5488 Value::Int(runmat_value::IntValue::U64(1_u64 << 63)),
5489 ],
5490 )
5491 .expect("queue transaction fill");
5492 call_builtin("DataTransaction.commit", &[tx]).expect("commit transaction");
5493
5494 let Value::Tensor(read_back) = call_builtin("DataArray.read", &[arr]).expect("read") else {
5495 panic!("expected tensor");
5496 };
5497 assert_eq!(
5498 read_back.integer_storage(),
5499 Some(&IntegerStorage::U64(vec![1_u64 << 63; 4]))
5500 );
5501 }
5502
5503 #[test]
5504 fn paired_complex_uint64_data_array_chunk_and_slice_paths_remain_exact() {
5505 let _serial = serial_test_guard();
5506 let _provider = native_provider_guard();
5507 let dir = tempfile::tempdir().expect("tempdir");
5508 let path = dir
5509 .path()
5510 .join("complex-uint64.data")
5511 .to_string_lossy()
5512 .to_string();
5513 let mut array_meta = StructValue::new();
5514 array_meta
5515 .fields
5516 .insert("dtype".to_string(), Value::String("uint64".to_string()));
5517 array_meta.fields.insert(
5518 "shape".to_string(),
5519 Value::Tensor(Tensor::new(vec![2.0, 2.0], vec![1, 2]).expect("shape")),
5520 );
5521 array_meta.fields.insert(
5522 "chunk".to_string(),
5523 Value::Tensor(Tensor::new(vec![1.0, 1.0], vec![1, 2]).expect("chunk")),
5524 );
5525 let mut arrays = StructValue::new();
5526 arrays
5527 .fields
5528 .insert("samples".to_string(), Value::Struct(array_meta));
5529 let mut schema = StructValue::new();
5530 schema
5531 .fields
5532 .insert("arrays".to_string(), Value::Struct(arrays));
5533 let ds = call_builtin("data.create", &[Value::String(path), Value::Struct(schema)])
5534 .expect("create dataset");
5535 let arr = call_builtin("Dataset.array", &[ds, Value::String("samples".to_string())])
5536 .expect("array");
5537 let wide = (1_u64 << 53) + 1;
5538 let initial_storage = IntegerComplexStorage::new(
5539 IntegerStorage::U64(vec![wide, u64::MAX, 3, 4]),
5540 IntegerStorage::U64(vec![u64::MAX, wide, 5, 6]),
5541 )
5542 .expect("initial paired storage");
5543 let initial = ComplexTensor::new_integer(initial_storage.clone(), vec![2, 2])
5544 .expect("initial paired tensor");
5545 call_builtin(
5546 "DataArray.write",
5547 &[arr.clone(), Value::ComplexTensor(initial)],
5548 )
5549 .expect("write paired array");
5550
5551 let Value::ComplexTensor(read_back) =
5552 call_builtin("DataArray.read", std::slice::from_ref(&arr)).expect("read paired array")
5553 else {
5554 panic!("expected paired tensor");
5555 };
5556 assert_eq!(read_back.integer_storage(), Some(&initial_storage));
5557
5558 let slice = Value::Cell(
5559 CellArray::new(
5560 vec![Value::Int(IntValue::I32(1)), Value::String(":".to_string())],
5561 1,
5562 2,
5563 )
5564 .expect("slice"),
5565 );
5566 let replacement_storage = IntegerComplexStorage::new(
5567 IntegerStorage::U64(vec![1_u64 << 63, u64::MAX - 1]),
5568 IntegerStorage::U64(vec![u64::MAX - 2, 1_u64 << 63]),
5569 )
5570 .expect("replacement storage");
5571 let replacement = ComplexTensor::new_integer(replacement_storage, vec![1, 2])
5572 .expect("replacement tensor");
5573 call_builtin(
5574 "DataArray.write",
5575 &[arr.clone(), slice, Value::ComplexTensor(replacement)],
5576 )
5577 .expect("write paired slice");
5578
5579 let Value::ComplexTensor(read_back) =
5580 call_builtin("DataArray.read", &[arr]).expect("read paired slice result")
5581 else {
5582 panic!("expected paired tensor");
5583 };
5584 let storage = read_back.integer_storage().expect("paired storage");
5585 assert_eq!(
5586 storage.real,
5587 IntegerStorage::U64(vec![1_u64 << 63, u64::MAX, u64::MAX - 1, 4])
5588 );
5589 assert_eq!(
5590 storage.imag,
5591 IntegerStorage::U64(vec![u64::MAX - 2, wide, 1_u64 << 63, 6])
5592 );
5593 }
5594}