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