1use std::cell::RefCell;
2use std::collections::{BTreeMap, BTreeSet};
3
4use serde::{Deserialize, Serialize};
5
6use crate::coordinator::{
7 parse_native_branch_view_filter, validate_fingerprint, CoordinatorDataPlanEnvelope,
8 CoordinatorRelation, CoordinatorRelationSet,
9};
10use crate::error::{DataError, Result};
11use crate::ids::{ObservationId, RepresentationId, SampleId, SourceId, TargetId};
12use crate::model::DataView;
13
14#[derive(Clone, Copy, Debug, Eq, PartialEq, Ord, PartialOrd, Serialize, Deserialize)]
15#[serde(rename_all = "snake_case")]
16pub enum CoordinatorHandleKind {
17 Data,
18 View,
19}
20
21#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
22pub struct CoordinatorHandleRef {
23 pub handle: u64,
24 pub kind: CoordinatorHandleKind,
25 pub owner_controller: String,
26}
27
28#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
29pub struct CoordinatorDataMaterializationRequest {
30 pub run_id: String,
31 pub node_id: String,
32 pub input_name: String,
33 pub phase: String,
34 #[serde(default)]
35 pub variant_id: Option<String>,
36 #[serde(default)]
37 pub fold_id: Option<String>,
38 pub request_id: String,
39 pub schema_fingerprint: String,
40 pub plan_fingerprint: String,
41 #[serde(default)]
42 pub relation_fingerprint: Option<String>,
43 pub output_representation: RepresentationId,
44 #[serde(default)]
45 pub source_ids: Vec<SourceId>,
46 #[serde(default)]
47 pub require_relations: bool,
48}
49
50impl CoordinatorDataMaterializationRequest {
51 pub fn validate(&self) -> Result<()> {
52 validate_non_empty("run_id", &self.run_id)?;
53 validate_non_empty("node_id", &self.node_id)?;
54 validate_non_empty("input_name", &self.input_name)?;
55 validate_non_empty("phase", &self.phase)?;
56 validate_non_empty("request_id", &self.request_id)?;
57 validate_fingerprint("schema", &self.schema_fingerprint)?;
58 validate_fingerprint("plan", &self.plan_fingerprint)?;
59 if let Some(relation_fingerprint) = &self.relation_fingerprint {
60 validate_fingerprint("relation", relation_fingerprint)?;
61 } else if self.require_relations {
62 return Err(DataError::Validation(format!(
63 "materialization request `{}` on `{}` requires relations but has no relation_fingerprint",
64 self.input_name, self.node_id
65 )));
66 }
67 let unique_sources = self.source_ids.iter().collect::<BTreeSet<_>>();
68 if unique_sources.len() != self.source_ids.len() {
69 return Err(DataError::Validation(format!(
70 "materialization request `{}` on `{}` contains duplicate source ids",
71 self.input_name, self.node_id
72 )));
73 }
74 Ok(())
75 }
76}
77
78#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
79pub struct CoordinatorDataHandleRecord {
80 pub handle: CoordinatorHandleRef,
81 pub run_id: String,
82 pub node_id: String,
83 pub input_name: String,
84 pub phase: String,
85 #[serde(default)]
86 pub variant_id: Option<String>,
87 #[serde(default)]
88 pub fold_id: Option<String>,
89 pub request_id: String,
90 pub schema_fingerprint: String,
91 pub plan_fingerprint: String,
92 #[serde(default)]
93 pub relation_fingerprint: Option<String>,
94 pub plan_id: String,
95 pub output_representation: RepresentationId,
96 #[serde(default)]
97 pub source_ids: Vec<SourceId>,
98 #[serde(default)]
99 pub sample_count: Option<usize>,
100 #[serde(default)]
101 pub relation_record_count: Option<usize>,
102}
103
104#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
105pub struct CoordinatorDataViewRecord {
106 pub handle: CoordinatorHandleRef,
107 pub parent_handle: CoordinatorHandleRef,
108 pub view: DataView,
109 pub sample_count: usize,
110 pub relation_record_count: usize,
111}
112
113#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
114pub struct CoordinatorTargetValue {
115 pub sample_id: SampleId,
116 pub value: serde_json::Value,
117}
118
119#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
120pub struct CoordinatorTargetTable {
121 pub target_id: TargetId,
122 pub values: Vec<CoordinatorTargetValue>,
123}
124
125impl CoordinatorTargetTable {
126 pub fn validate(&self) -> Result<()> {
127 if self.values.is_empty() {
128 return Err(DataError::Validation(format!(
129 "target table `{}` contains no values",
130 self.target_id
131 )));
132 }
133 let mut seen = BTreeSet::new();
134 for value in &self.values {
135 if !seen.insert(&value.sample_id) {
136 return Err(DataError::Validation(format!(
137 "target table `{}` contains duplicate sample `{}`",
138 self.target_id, value.sample_id
139 )));
140 }
141 }
142 Ok(())
143 }
144}
145
146#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
147pub struct CoordinatorTargetBlock {
148 pub target_id: TargetId,
149 pub sample_ids: Vec<SampleId>,
150 pub values: Vec<serde_json::Value>,
151}
152
153#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
154pub struct CoordinatorMultiTargetBlock {
155 pub target_ids: Vec<TargetId>,
156 pub sample_ids: Vec<SampleId>,
157 pub values: Vec<Vec<serde_json::Value>>,
159 pub validity_masks: Vec<Vec<bool>>,
161}
162
163#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
164pub struct CoordinatorFeatureRow {
165 pub observation_id: ObservationId,
166 pub values: Vec<serde_json::Value>,
167}
168
169#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
170pub struct CoordinatorFeatureTable {
171 pub feature_set_id: String,
172 pub representation_id: RepresentationId,
173 pub feature_names: Vec<String>,
174 pub rows: Vec<CoordinatorFeatureRow>,
175}
176
177impl CoordinatorFeatureTable {
178 pub fn validate(&self) -> Result<()> {
179 validate_non_empty("feature_set_id", &self.feature_set_id)?;
180 if self.feature_names.is_empty() {
181 return Err(DataError::Validation(format!(
182 "feature table `{}` contains no features",
183 self.feature_set_id
184 )));
185 }
186 let mut seen_features = BTreeSet::new();
187 for feature_name in &self.feature_names {
188 validate_non_empty("feature_name", feature_name)?;
189 if !seen_features.insert(feature_name) {
190 return Err(DataError::Validation(format!(
191 "feature table `{}` contains duplicate feature `{}`",
192 self.feature_set_id, feature_name
193 )));
194 }
195 }
196 if self.rows.is_empty() {
197 return Err(DataError::Validation(format!(
198 "feature table `{}` contains no rows",
199 self.feature_set_id
200 )));
201 }
202 let mut seen_observations = BTreeSet::new();
203 for row in &self.rows {
204 if !seen_observations.insert(&row.observation_id) {
205 return Err(DataError::Validation(format!(
206 "feature table `{}` contains duplicate observation `{}`",
207 self.feature_set_id, row.observation_id
208 )));
209 }
210 if row.values.len() != self.feature_names.len() {
211 return Err(DataError::Validation(format!(
212 "feature table `{}` row `{}` has {} values for {} features",
213 self.feature_set_id,
214 row.observation_id,
215 row.values.len(),
216 self.feature_names.len()
217 )));
218 }
219 }
220 Ok(())
221 }
222}
223
224#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
225pub struct CoordinatorFeatureBlock {
226 pub feature_set_id: String,
227 pub representation_id: RepresentationId,
228 pub feature_names: Vec<String>,
229 pub observation_ids: Vec<ObservationId>,
230 pub sample_ids: Vec<SampleId>,
231 pub values: Vec<Vec<serde_json::Value>>,
232}
233
234#[derive(Clone, Debug, PartialEq)]
240pub struct CoordinatorFeatureBlockF64 {
241 pub feature_set_id: String,
242 pub representation_id: RepresentationId,
243 pub feature_names: Vec<String>,
244 pub observation_ids: Vec<ObservationId>,
245 pub sample_ids: Vec<SampleId>,
246 pub values: Vec<f64>,
247}
248
249#[derive(Debug)]
250pub struct CoordinatorHandleArena {
251 owner_controller: String,
252 next_handle: RefCell<u64>,
253 records: RefCell<BTreeMap<u64, CoordinatorDataHandleRecord>>,
254 data_relations: RefCell<BTreeMap<u64, CoordinatorRelationSet>>,
255 view_records: RefCell<BTreeMap<u64, CoordinatorDataViewRecord>>,
256 view_relations: RefCell<BTreeMap<u64, CoordinatorRelationSet>>,
257}
258
259impl CoordinatorHandleArena {
260 pub fn new(owner_controller: impl Into<String>) -> Result<Self> {
261 let owner_controller = owner_controller.into();
262 validate_non_empty("owner_controller", &owner_controller)?;
263 Ok(Self {
264 owner_controller,
265 next_handle: RefCell::new(1),
266 records: RefCell::new(BTreeMap::new()),
267 data_relations: RefCell::new(BTreeMap::new()),
268 view_records: RefCell::new(BTreeMap::new()),
269 view_relations: RefCell::new(BTreeMap::new()),
270 })
271 }
272
273 pub fn materialize(
274 &self,
275 envelope: &CoordinatorDataPlanEnvelope,
276 request: &CoordinatorDataMaterializationRequest,
277 ) -> Result<CoordinatorDataHandleRecord> {
278 envelope.validate()?;
279 request.validate()?;
280 validate_request_against_envelope(envelope, request)?;
281 let scoped_relations = envelope
282 .coordinator_relations
283 .as_ref()
284 .map(|relations| scoped_relations_for_materialization(relations, request))
285 .transpose()?;
286
287 let handle = CoordinatorHandleRef {
288 handle: self.next_handle(),
289 kind: CoordinatorHandleKind::Data,
290 owner_controller: self.owner_controller.clone(),
291 };
292 let record = CoordinatorDataHandleRecord {
293 handle: handle.clone(),
294 run_id: request.run_id.clone(),
295 node_id: request.node_id.clone(),
296 input_name: request.input_name.clone(),
297 phase: request.phase.clone(),
298 variant_id: request.variant_id.clone(),
299 fold_id: request.fold_id.clone(),
300 request_id: request.request_id.clone(),
301 schema_fingerprint: request.schema_fingerprint.clone(),
302 plan_fingerprint: request.plan_fingerprint.clone(),
303 relation_fingerprint: request.relation_fingerprint.clone(),
304 plan_id: envelope.plan.id.clone(),
305 output_representation: request.output_representation.clone(),
306 source_ids: request.source_ids.clone(),
307 sample_count: scoped_relations.as_ref().map(|relations| {
308 relations
309 .records
310 .iter()
311 .map(|record| &record.sample_id)
312 .collect::<BTreeSet<_>>()
313 .len()
314 }),
315 relation_record_count: scoped_relations
316 .as_ref()
317 .map(|relations| relations.records.len()),
318 };
319 self.records
320 .borrow_mut()
321 .insert(handle.handle, record.clone());
322 if let Some(relations) = scoped_relations {
323 self.data_relations
324 .borrow_mut()
325 .insert(handle.handle, relations);
326 }
327 Ok(record)
328 }
329
330 pub fn make_view(
331 &self,
332 data_handle: u64,
333 view: &DataView,
334 ) -> Result<CoordinatorDataViewRecord> {
335 validate_view(view)?;
336 let parent =
337 self.records
338 .borrow()
339 .get(&data_handle)
340 .cloned()
341 .ok_or(DataError::UnknownHandle {
342 kind: "data",
343 handle: data_handle,
344 })?;
345 let relations = self
346 .data_relations
347 .borrow()
348 .get(&data_handle)
349 .cloned()
350 .ok_or_else(|| {
351 DataError::Validation(format!(
352 "data handle `{data_handle}` has no coordinator relations"
353 ))
354 })?;
355 let filtered = filter_relations(&relations.records, view)?;
356 let sample_count = unique_sample_count(&filtered);
357 let relation_record_count = filtered.len();
358 let handle = CoordinatorHandleRef {
359 handle: self.next_handle(),
360 kind: CoordinatorHandleKind::View,
361 owner_controller: self.owner_controller.clone(),
362 };
363 let record = CoordinatorDataViewRecord {
364 handle: handle.clone(),
365 parent_handle: parent.handle,
366 view: view.clone(),
367 sample_count,
368 relation_record_count,
369 };
370 self.view_records
371 .borrow_mut()
372 .insert(handle.handle, record.clone());
373 self.view_relations
374 .borrow_mut()
375 .insert(handle.handle, CoordinatorRelationSet { records: filtered });
376 Ok(record)
377 }
378
379 pub fn view_record(&self, handle: u64) -> Option<CoordinatorDataViewRecord> {
380 self.view_records.borrow().get(&handle).cloned()
381 }
382
383 pub fn view_identity(&self, handle: u64) -> Result<CoordinatorRelationSet> {
384 self.view_relations
385 .borrow()
386 .get(&handle)
387 .cloned()
388 .ok_or(DataError::UnknownHandle {
389 kind: "view",
390 handle,
391 })
392 }
393
394 pub fn data_identity(&self, handle: u64) -> Result<CoordinatorRelationSet> {
395 if !self.records.borrow().contains_key(&handle) {
400 return Err(DataError::UnknownHandle {
401 kind: "data",
402 handle,
403 });
404 }
405 self.data_relations
406 .borrow()
407 .get(&handle)
408 .cloned()
409 .ok_or_else(|| {
410 DataError::Validation(format!(
411 "data handle `{handle}` has no coordinator relations"
412 ))
413 })
414 }
415
416 pub fn release_handle(&self, handle: u64) -> bool {
417 if self.view_records.borrow_mut().remove(&handle).is_some() {
418 self.view_relations.borrow_mut().remove(&handle);
419 return true;
420 }
421 if let Some(record) = self.records.borrow_mut().remove(&handle) {
422 self.data_relations.borrow_mut().remove(&handle);
423 let child_views = self
424 .view_records
425 .borrow()
426 .iter()
427 .filter_map(|(view_handle, view_record)| {
428 (view_record.parent_handle == record.handle).then_some(*view_handle)
429 })
430 .collect::<Vec<_>>();
431 for view_handle in child_views {
432 self.view_records.borrow_mut().remove(&view_handle);
433 self.view_relations.borrow_mut().remove(&view_handle);
434 }
435 return true;
436 }
437 false
438 }
439
440 pub fn target_values(
441 &self,
442 view_handle: u64,
443 target_table: &CoordinatorTargetTable,
444 ) -> Result<CoordinatorTargetBlock> {
445 target_table.validate()?;
446 let relations = self.view_identity(view_handle)?;
447 let values_by_sample = target_table
448 .values
449 .iter()
450 .map(|value| (&value.sample_id, &value.value))
451 .collect::<BTreeMap<_, _>>();
452 let mut seen_samples = BTreeSet::new();
453 let mut sample_ids = Vec::new();
454 let mut values = Vec::new();
455 for relation in relations.records.iter().filter(|relation| {
456 relation
457 .target_id
458 .as_ref()
459 .map(|target_id| target_id == &target_table.target_id)
460 .unwrap_or(true)
461 }) {
462 if !seen_samples.insert(&relation.sample_id) {
463 continue;
464 }
465 let value = values_by_sample.get(&relation.sample_id).ok_or_else(|| {
466 DataError::Validation(format!(
467 "target table `{}` has no value for sample `{}`",
468 target_table.target_id, relation.sample_id
469 ))
470 })?;
471 sample_ids.push(relation.sample_id.clone());
472 values.push((*value).clone());
473 }
474 if sample_ids.is_empty() {
475 return Err(DataError::Validation(format!(
476 "view `{view_handle}` contains no samples for target `{}`",
477 target_table.target_id
478 )));
479 }
480 Ok(CoordinatorTargetBlock {
481 target_id: target_table.target_id.clone(),
482 sample_ids,
483 values,
484 })
485 }
486
487 pub fn multi_target_values(
488 &self,
489 view_handle: u64,
490 target_tables: &[CoordinatorTargetTable],
491 ) -> Result<CoordinatorMultiTargetBlock> {
492 if target_tables.is_empty() {
493 return Err(DataError::Validation(
494 "multi-target materialization requires at least one target table".to_string(),
495 ));
496 }
497 let mut seen_targets = BTreeSet::new();
498 for table in target_tables {
499 table.validate()?;
500 if !seen_targets.insert(table.target_id.clone()) {
501 return Err(DataError::Validation(format!(
502 "multi-target materialization contains duplicate target `{}`",
503 table.target_id
504 )));
505 }
506 }
507
508 let target_ids = target_tables
509 .iter()
510 .map(|table| table.target_id.clone())
511 .collect::<Vec<_>>();
512 let target_universe = target_ids.iter().collect::<BTreeSet<_>>();
513 let relations = self.view_identity(view_handle)?;
514 let mut seen_samples = BTreeSet::new();
515 let mut sample_ids = Vec::new();
516 for relation in relations.records.iter().filter(|relation| {
517 relation
518 .target_id
519 .as_ref()
520 .map(|target_id| target_universe.contains(target_id))
521 .unwrap_or(true)
522 }) {
523 if seen_samples.insert(relation.sample_id.clone()) {
524 sample_ids.push(relation.sample_id.clone());
525 }
526 }
527 if sample_ids.is_empty() {
528 return Err(DataError::Validation(format!(
529 "view `{view_handle}` contains no samples for requested targets"
530 )));
531 }
532
533 let mut values = Vec::with_capacity(target_tables.len());
534 let mut validity_masks = Vec::with_capacity(target_tables.len());
535 for table in target_tables {
536 let values_by_sample = table
537 .values
538 .iter()
539 .map(|value| (&value.sample_id, &value.value))
540 .collect::<BTreeMap<_, _>>();
541 let mut target_values = Vec::with_capacity(sample_ids.len());
542 let mut validity = Vec::with_capacity(sample_ids.len());
543 for sample_id in &sample_ids {
544 match values_by_sample.get(sample_id) {
545 Some(value) if !value.is_null() => {
546 target_values.push((*value).clone());
547 validity.push(true);
548 }
549 Some(_) | None => {
550 target_values.push(serde_json::Value::Null);
551 validity.push(false);
552 }
553 }
554 }
555 values.push(target_values);
556 validity_masks.push(validity);
557 }
558
559 Ok(CoordinatorMultiTargetBlock {
560 target_ids,
561 sample_ids,
562 values,
563 validity_masks,
564 })
565 }
566
567 pub fn feature_values(
568 &self,
569 view_handle: u64,
570 feature_table: &CoordinatorFeatureTable,
571 ) -> Result<CoordinatorFeatureBlock> {
572 feature_table.validate()?;
573 let view_record = self
574 .view_records
575 .borrow()
576 .get(&view_handle)
577 .cloned()
578 .ok_or(DataError::UnknownHandle {
579 kind: "view",
580 handle: view_handle,
581 })?;
582 let parent_record = self
583 .records
584 .borrow()
585 .get(&view_record.parent_handle.handle)
586 .cloned()
587 .ok_or(DataError::UnknownHandle {
588 kind: "data",
589 handle: view_record.parent_handle.handle,
590 })?;
591 if feature_table.representation_id != parent_record.output_representation {
592 return Err(DataError::Validation(format!(
593 "feature table `{}` representation `{}` does not match materialized output representation `{}`",
594 feature_table.feature_set_id,
595 feature_table.representation_id,
596 parent_record.output_representation
597 )));
598 }
599 let relations = self.view_identity(view_handle)?;
600 let selected_indices = selected_feature_indices(feature_table, &view_record.view)?;
601 let rows_by_observation = feature_table
602 .rows
603 .iter()
604 .map(|row| (&row.observation_id, row))
605 .collect::<BTreeMap<_, _>>();
606 let mut observation_ids = Vec::with_capacity(relations.records.len());
607 let mut sample_ids = Vec::with_capacity(relations.records.len());
608 let mut values = Vec::with_capacity(relations.records.len());
609 for relation in &relations.records {
610 let row = rows_by_observation
611 .get(&relation.observation_id)
612 .ok_or_else(|| {
613 DataError::Validation(format!(
614 "feature table `{}` has no row for observation `{}`",
615 feature_table.feature_set_id, relation.observation_id
616 ))
617 })?;
618 observation_ids.push(relation.observation_id.clone());
619 sample_ids.push(relation.sample_id.clone());
620 values.push(
621 selected_indices
622 .iter()
623 .map(|idx| row.values[*idx].clone())
624 .collect(),
625 );
626 }
627 Ok(CoordinatorFeatureBlock {
628 feature_set_id: feature_table.feature_set_id.clone(),
629 representation_id: feature_table.representation_id.clone(),
630 feature_names: selected_indices
631 .iter()
632 .map(|idx| feature_table.feature_names[*idx].clone())
633 .collect(),
634 observation_ids,
635 sample_ids,
636 values,
637 })
638 }
639
640 pub fn handle_record(&self, handle: u64) -> Option<CoordinatorDataHandleRecord> {
641 self.records.borrow().get(&handle).cloned()
642 }
643
644 pub fn handle_records(&self) -> Vec<CoordinatorDataHandleRecord> {
645 self.records.borrow().values().cloned().collect()
646 }
647
648 fn next_handle(&self) -> u64 {
649 let mut next = self.next_handle.borrow_mut();
650 let handle = *next;
651 *next += 1;
652 handle
653 }
654}
655
656fn validate_request_against_envelope(
657 envelope: &CoordinatorDataPlanEnvelope,
658 request: &CoordinatorDataMaterializationRequest,
659) -> Result<()> {
660 if request.schema_fingerprint != envelope.schema_fingerprint {
661 return Err(DataError::FingerprintMismatch {
662 kind: "schema",
663 expected: envelope.schema_fingerprint.clone(),
664 actual: request.schema_fingerprint.clone(),
665 });
666 }
667 if request.plan_fingerprint != envelope.plan_fingerprint {
668 return Err(DataError::FingerprintMismatch {
669 kind: "plan",
670 expected: envelope.plan_fingerprint.clone(),
671 actual: request.plan_fingerprint.clone(),
672 });
673 }
674 if request.relation_fingerprint != envelope.relation_fingerprint {
675 let none = || "<none>".to_string();
676 return Err(DataError::FingerprintMismatch {
677 kind: "relation",
678 expected: envelope.relation_fingerprint.clone().unwrap_or_else(none),
679 actual: request.relation_fingerprint.clone().unwrap_or_else(none),
680 });
681 }
682 if request.require_relations && envelope.coordinator_relations.is_none() {
683 return Err(DataError::Validation(format!(
684 "materialization request `{}` on `{}` requires coordinator relations",
685 request.input_name, request.node_id
686 )));
687 }
688 if request.output_representation != envelope.plan.output_representation {
689 return Err(DataError::Validation(format!(
690 "materialization request `{}` on `{}` output representation `{}` does not match plan output `{}`",
691 request.input_name,
692 request.node_id,
693 request.output_representation,
694 envelope.plan.output_representation
695 )));
696 }
697 if !request.source_ids.is_empty() {
698 let plan_sources = envelope
699 .plan
700 .steps
701 .iter()
702 .filter_map(|step| step.source_id.as_ref())
703 .collect::<BTreeSet<_>>();
704 for source_id in &request.source_ids {
705 if !plan_sources.contains(source_id) {
706 return Err(DataError::Validation(format!(
707 "materialization request `{}` on `{}` source `{}` is not present in data plan `{}`",
708 request.input_name, request.node_id, source_id, envelope.plan.id
709 )));
710 }
711 }
712 }
713 Ok(())
714}
715
716fn scoped_relations_for_materialization(
717 relations: &CoordinatorRelationSet,
718 request: &CoordinatorDataMaterializationRequest,
719) -> Result<CoordinatorRelationSet> {
720 relations.validate()?;
721 if request.source_ids.is_empty() {
722 return Ok(relations.clone());
723 }
724 let source_filter = request.source_ids.iter().collect::<BTreeSet<_>>();
725 let scoped = relations
726 .records
727 .iter()
728 .filter(|relation| {
729 relation
730 .source_id
731 .as_ref()
732 .map(|source_id| source_filter.contains(source_id))
733 .unwrap_or(false)
734 })
735 .cloned()
736 .collect::<Vec<_>>();
737 if scoped.is_empty() {
738 return Err(DataError::Validation(format!(
739 "materialization request `{}` on `{}` selected no coordinator relations for requested source ids",
740 request.input_name, request.node_id
741 )));
742 }
743 let scoped = CoordinatorRelationSet { records: scoped };
744 scoped.validate()?;
745 Ok(scoped)
746}
747
748fn validate_non_empty(label: &str, value: &str) -> Result<()> {
749 if value.trim().is_empty() {
750 return Err(DataError::Validation(format!("{label} must not be empty")));
751 }
752 Ok(())
753}
754
755fn validate_view(view: &DataView) -> Result<()> {
756 if let Some(samples) = &view.sample_ids {
757 let unique = samples.iter().collect::<BTreeSet<_>>();
758 if unique.len() != samples.len() {
759 return Err(DataError::Validation(
760 "data view contains duplicate sample ids".to_string(),
761 ));
762 }
763 }
764 if let Some(sources) = &view.source_ids {
765 let unique = sources.iter().collect::<BTreeSet<_>>();
766 if unique.len() != sources.len() {
767 return Err(DataError::Validation(
768 "data view contains duplicate source ids".to_string(),
769 ));
770 }
771 }
772 if let Some(columns) = &view.columns {
773 let unique = columns.iter().collect::<BTreeSet<_>>();
774 if unique.len() != columns.len() {
775 return Err(DataError::Validation(
776 "data view contains duplicate columns".to_string(),
777 ));
778 }
779 if columns.iter().any(|column| column.trim().is_empty()) {
780 return Err(DataError::Validation(
781 "data view contains an empty column".to_string(),
782 ));
783 }
784 }
785 if let Some(branch_view) = &view.branch_view {
786 branch_view.validate()?;
787 }
788 Ok(())
789}
790
791fn filter_relations(
792 relations: &[CoordinatorRelation],
793 view: &DataView,
794) -> Result<Vec<CoordinatorRelation>> {
795 let sample_filter = view
796 .sample_ids
797 .as_ref()
798 .map(|sample_ids| sample_ids.iter().collect::<BTreeSet<_>>());
799 let source_filter = view
800 .source_ids
801 .as_ref()
802 .map(|source_ids| source_ids.iter().collect::<BTreeSet<_>>());
803 let mut branch_source_filter: Option<BTreeSet<&SourceId>> = None;
810 let mut branch_metadata_filter: Option<&BTreeMap<String, serde_json::Value>> = None;
811 let mut branch_tag_filter: Option<BTreeSet<&String>> = None;
812 let mut branch_native_filter = None;
813 if let Some(branch_view) = view.branch_view.as_ref() {
814 match branch_view.mode {
815 crate::coordinator::CoordinatorBranchViewMode::BySource => {
816 branch_source_filter = Some(
817 branch_view
818 .selector
819 .source_ids
820 .iter()
821 .collect::<BTreeSet<_>>(),
822 );
823 }
824 crate::coordinator::CoordinatorBranchViewMode::ByMetadata => {
828 branch_metadata_filter = Some(&branch_view.selector.metadata);
829 }
830 crate::coordinator::CoordinatorBranchViewMode::ByTag => {
831 branch_tag_filter = Some(branch_view.selector.tags.iter().collect::<BTreeSet<_>>());
832 }
833 crate::coordinator::CoordinatorBranchViewMode::Separation => {}
834 crate::coordinator::CoordinatorBranchViewMode::ByFilter => {
835 let filter = branch_view.selector.filter.as_ref().ok_or_else(|| {
836 DataError::Validation(format!(
837 "coordinator branch view `{}` mode=by_filter requires filter",
838 branch_view.view_id
839 ))
840 })?;
841 branch_native_filter = Some(parse_native_branch_view_filter(
842 filter,
843 &format!("coordinator branch view `{}`", branch_view.view_id),
844 )?);
845 }
846 }
847 }
848 let mut filtered = relations
849 .iter()
850 .enumerate()
851 .filter(|relation| {
852 let relation = relation.1;
853 sample_filter
854 .as_ref()
855 .map(|samples| samples.contains(&relation.sample_id))
856 .unwrap_or(true)
857 })
858 .filter(|relation| {
859 let relation = relation.1;
860 source_filter
861 .as_ref()
862 .map(|sources| {
863 relation
864 .source_id
865 .as_ref()
866 .map(|source_id| sources.contains(source_id))
867 .unwrap_or(false)
868 })
869 .unwrap_or(true)
870 })
871 .filter(|relation| {
872 let relation = relation.1;
873 branch_source_filter
874 .as_ref()
875 .map(|sources| {
876 relation
877 .source_id
878 .as_ref()
879 .map(|source_id| sources.contains(source_id))
880 .unwrap_or(false)
881 })
882 .unwrap_or(true)
883 })
884 .filter(|relation| {
885 let relation = relation.1;
886 branch_metadata_filter
887 .map(|selector| {
888 selector
889 .iter()
890 .all(|(key, value)| relation.metadata.get(key) == Some(value))
891 })
892 .unwrap_or(true)
893 })
894 .filter(|relation| {
895 let relation = relation.1;
896 branch_tag_filter
897 .as_ref()
898 .map(|tags| tags.iter().all(|tag| relation.tags.contains(tag)))
899 .unwrap_or(true)
900 })
901 .filter(|relation| {
902 let relation = relation.1;
903 branch_native_filter
904 .as_ref()
905 .map(|filter| {
906 filter
907 .metadata_equals
908 .iter()
909 .all(|(key, value)| relation.metadata.get(key) == Some(value))
910 && filter
911 .tags_all
912 .iter()
913 .all(|tag| relation.tags.contains(tag))
914 })
915 .unwrap_or(true)
916 })
917 .filter(|relation| view.include_augmented || !relation.1.is_augmented)
918 .filter(|relation| view.include_excluded || !relation.1.excluded)
919 .map(|(idx, relation)| (idx, relation.clone()))
920 .collect::<Vec<_>>();
921 if filtered.is_empty() {
922 return Err(DataError::Validation(
923 "data view selected no coordinator relations".to_string(),
924 ));
925 }
926 if let Some(sample_ids) = &view.sample_ids {
927 let sample_order = sample_ids
928 .iter()
929 .enumerate()
930 .map(|(idx, sample_id)| (sample_id, idx))
931 .collect::<BTreeMap<_, _>>();
932 filtered.sort_by_key(|(idx, relation)| {
933 (
934 sample_order
935 .get(&relation.sample_id)
936 .copied()
937 .unwrap_or(usize::MAX),
938 *idx,
939 )
940 });
941 }
942 Ok(filtered.into_iter().map(|(_, relation)| relation).collect())
943}
944
945fn unique_sample_count(relations: &[CoordinatorRelation]) -> usize {
946 relations
947 .iter()
948 .map(|relation| &relation.sample_id)
949 .collect::<BTreeSet<_>>()
950 .len()
951}
952
953fn selected_feature_indices(
954 table: &CoordinatorFeatureTable,
955 view: &DataView,
956) -> Result<Vec<usize>> {
957 let index_by_name = table
958 .feature_names
959 .iter()
960 .enumerate()
961 .map(|(idx, name)| (name, idx))
962 .collect::<BTreeMap<_, _>>();
963 let indices = if let Some(columns) = &view.columns {
964 columns
965 .iter()
966 .map(|column| {
967 index_by_name.get(column).copied().ok_or_else(|| {
968 DataError::Validation(format!(
969 "feature table `{}` has no feature column `{}`",
970 table.feature_set_id, column
971 ))
972 })
973 })
974 .collect::<Result<Vec<_>>>()?
975 } else {
976 (0..table.feature_names.len()).collect()
977 };
978 if indices.is_empty() {
979 return Err(DataError::Validation(format!(
980 "feature table `{}` selected no feature columns",
981 table.feature_set_id
982 )));
983 }
984 Ok(indices)
985}
986
987#[cfg(test)]
988mod tests {
989 use super::*;
990 use serde_json::json;
991
992 fn envelope() -> CoordinatorDataPlanEnvelope {
993 serde_json::from_str(include_str!(
994 "../../../examples/fixtures/oof_campaign/coordinator_data_plan_envelope_nir.json"
995 ))
996 .unwrap()
997 }
998
999 fn request() -> CoordinatorDataMaterializationRequest {
1000 serde_json::from_str(include_str!(
1001 "../../../examples/fixtures/oof_campaign/materialization_request_model_base_x.json"
1002 ))
1003 .unwrap()
1004 }
1005
1006 #[test]
1007 fn materializes_validated_coordinator_handle_record() {
1008 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1009 let record = arena.materialize(&envelope(), &request()).unwrap();
1010
1011 assert_eq!(record.handle.handle, 1);
1012 assert_eq!(record.handle.kind, CoordinatorHandleKind::Data);
1013 assert_eq!(record.input_name, "x");
1014 assert_eq!(record.plan_id, "nir-to-tabular");
1015 assert_eq!(record.sample_count, Some(2));
1016 assert_eq!(record.relation_record_count, Some(4));
1017 assert_eq!(arena.handle_record(1), Some(record));
1018 assert_eq!(arena.handle_records().len(), 1);
1019 }
1020
1021 #[test]
1022 fn materialization_refuses_fingerprint_mismatch() {
1023 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1024 let mut request = request();
1025 request.plan_fingerprint = "0".repeat(64);
1026
1027 assert!(arena.materialize(&envelope(), &request).is_err());
1028 }
1029
1030 #[test]
1031 fn materialization_scopes_relations_to_requested_sources() {
1032 let mut envelope = envelope();
1033 let chem = SourceId::new("chem").unwrap();
1034 envelope.plan.steps.push(crate::plan::DataPlanStep {
1035 kind: crate::plan::DataPlanStepKind::Materialize,
1036 source_id: Some(chem.clone()),
1037 adapter_id: None,
1038 input_representation: None,
1039 output_representation: Some(RepresentationId::new("tabular_numeric").unwrap()),
1040 fit_scope: crate::plan::FitScope::Stateless,
1041 requires_user_choice: false,
1042 metadata: BTreeMap::new(),
1043 });
1044 envelope.plan_fingerprint = crate::data_plan_fingerprint(&envelope.plan).unwrap();
1045 envelope.relation_fingerprint = None;
1046 envelope
1047 .coordinator_relations
1048 .as_mut()
1049 .unwrap()
1050 .records
1051 .push(CoordinatorRelation {
1052 observation_id: ObservationId::new("chem.S001").unwrap(),
1053 sample_id: SampleId::new("S001").unwrap(),
1054 target_id: Some(TargetId::new("y").unwrap()),
1055 group_id: None,
1056 origin_sample_id: None,
1057 source_id: Some(chem.clone()),
1058 is_augmented: false,
1059 excluded: false,
1060 metadata: BTreeMap::new(),
1061 tags: Vec::new(),
1062 });
1063 envelope.validate().unwrap();
1064
1065 let mut request = request();
1066 request.plan_fingerprint = envelope.plan_fingerprint.clone();
1067 request.relation_fingerprint = None;
1068 request.require_relations = false;
1069 request.source_ids = vec![SourceId::new("nir").unwrap()];
1070
1071 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1072 let data = arena.materialize(&envelope, &request).unwrap();
1073 let view = arena
1074 .make_view(data.handle.handle, &DataView::default())
1075 .unwrap();
1076 let identity = arena.view_identity(view.handle.handle).unwrap();
1077
1078 assert_eq!(data.relation_record_count, Some(4));
1079 assert_eq!(
1080 arena
1081 .data_identity(data.handle.handle)
1082 .unwrap()
1083 .records
1084 .len(),
1085 4
1086 );
1087 assert!(identity
1088 .records
1089 .iter()
1090 .all(|record| record.source_id.as_ref() == Some(&SourceId::new("nir").unwrap())));
1091 }
1092
1093 #[test]
1094 fn view_filters_augmented_rows_and_preserves_repetition_identity() {
1095 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1096 let data = arena.materialize(&envelope(), &request()).unwrap();
1097 let view = DataView {
1098 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1099 include_augmented: false,
1100 ..Default::default()
1101 };
1102
1103 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1104 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1105
1106 assert_eq!(view_record.handle.kind, CoordinatorHandleKind::View);
1107 assert_eq!(view_record.sample_count, 1);
1108 assert_eq!(view_record.relation_record_count, 2);
1109 assert_eq!(identity.records.len(), 2);
1110 assert_eq!(identity.records[0].observation_id.as_str(), "obs.S001.base");
1111 assert_eq!(identity.records[1].observation_id.as_str(), "obs.S001.rep1");
1112 assert_eq!(
1113 arena.view_record(view_record.handle.handle),
1114 Some(view_record)
1115 );
1116 }
1117
1118 #[test]
1119 fn view_drops_excluded_rows_for_training_and_keeps_them_otherwise() {
1120 let mut envelope = envelope();
1124 for record in &mut envelope.coordinator_relations.as_mut().unwrap().records {
1125 if record.observation_id.as_str() == "obs.S002.base" {
1126 record.excluded = true;
1127 }
1128 }
1129 envelope.validate().unwrap();
1134
1135 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1136 let data = arena.materialize(&envelope, &request()).unwrap();
1137
1138 let train_view = DataView {
1140 include_augmented: true,
1141 include_excluded: false,
1142 ..Default::default()
1143 };
1144 let train_record = arena.make_view(data.handle.handle, &train_view).unwrap();
1145 let train_identity = arena.view_identity(train_record.handle.handle).unwrap();
1146 assert!(
1147 train_identity
1148 .records
1149 .iter()
1150 .all(|record| record.sample_id.as_str() != "S002"),
1151 "excluded sample S002 must be absent from a training view"
1152 );
1153 assert_eq!(train_record.relation_record_count, 3);
1154
1155 let predict_view = DataView {
1157 include_augmented: true,
1158 include_excluded: true,
1159 ..Default::default()
1160 };
1161 let predict_record = arena.make_view(data.handle.handle, &predict_view).unwrap();
1162 let predict_identity = arena.view_identity(predict_record.handle.handle).unwrap();
1163 assert!(
1164 predict_identity
1165 .records
1166 .iter()
1167 .any(|record| record.sample_id.as_str() == "S002"),
1168 "excluded sample S002 must be present in a validation/predict view"
1169 );
1170 assert_eq!(predict_record.relation_record_count, 4);
1171 }
1172
1173 #[test]
1174 fn branch_view_by_source_filters_relations_to_branch_sources() {
1175 use crate::coordinator::{
1176 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1177 };
1178
1179 let mut envelope = envelope();
1180 let chem = SourceId::new("chem").unwrap();
1181 envelope.plan.steps.push(crate::plan::DataPlanStep {
1182 kind: crate::plan::DataPlanStepKind::Materialize,
1183 source_id: Some(chem.clone()),
1184 adapter_id: None,
1185 input_representation: None,
1186 output_representation: Some(RepresentationId::new("tabular_numeric").unwrap()),
1187 fit_scope: crate::plan::FitScope::Stateless,
1188 requires_user_choice: false,
1189 metadata: BTreeMap::new(),
1190 });
1191 envelope.plan_fingerprint = crate::data_plan_fingerprint(&envelope.plan).unwrap();
1192 envelope.relation_fingerprint = None;
1193 envelope
1194 .coordinator_relations
1195 .as_mut()
1196 .unwrap()
1197 .records
1198 .push(CoordinatorRelation {
1199 observation_id: ObservationId::new("chem.S001").unwrap(),
1200 sample_id: SampleId::new("S001").unwrap(),
1201 target_id: Some(TargetId::new("y").unwrap()),
1202 group_id: None,
1203 origin_sample_id: None,
1204 source_id: Some(chem.clone()),
1205 is_augmented: false,
1206 excluded: false,
1207 metadata: BTreeMap::new(),
1208 tags: Vec::new(),
1209 });
1210 envelope.validate().unwrap();
1211 let mut request = request();
1212 request.plan_fingerprint = envelope.plan_fingerprint.clone();
1213 request.relation_fingerprint = None;
1214 request.require_relations = false;
1215 request.source_ids = vec![SourceId::new("nir").unwrap(), chem.clone()];
1216
1217 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1218 let data = arena.materialize(&envelope, &request).unwrap();
1219 let view = DataView {
1220 branch_view: Some(CoordinatorBranchView {
1221 view_id: "branch_view:nir".to_string(),
1222 branch_id: "branch:nir_only".to_string(),
1223 mode: CoordinatorBranchViewMode::BySource,
1224 selector: CoordinatorBranchViewSelector {
1225 source_ids: vec![SourceId::new("nir").unwrap()],
1226 ..Default::default()
1227 },
1228 allow_overlap: false,
1229 metadata: BTreeMap::new(),
1230 }),
1231 include_augmented: true,
1232 ..Default::default()
1233 };
1234 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1235 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1236
1237 assert!(identity
1238 .records
1239 .iter()
1240 .all(|record| record.source_id.as_ref() == Some(&SourceId::new("nir").unwrap())));
1241 assert!(identity.records.iter().any(|record| record
1242 .observation_id
1243 .as_str()
1244 .starts_with("obs.S001")
1245 || record.observation_id.as_str().starts_with("obs.S002")));
1246 }
1247
1248 #[test]
1249 fn branch_view_separation_does_not_restrict_relations() {
1250 use crate::coordinator::{
1251 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1252 };
1253
1254 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1255 let data = arena.materialize(&envelope(), &request()).unwrap();
1256 let view = DataView {
1257 branch_view: Some(CoordinatorBranchView {
1258 view_id: "branch_view:separation".to_string(),
1259 branch_id: "branch:0".to_string(),
1260 mode: CoordinatorBranchViewMode::Separation,
1261 selector: CoordinatorBranchViewSelector {
1262 tags: vec!["clean".to_string()],
1263 ..Default::default()
1264 },
1265 allow_overlap: false,
1266 metadata: BTreeMap::new(),
1267 }),
1268 include_augmented: true,
1269 ..Default::default()
1270 };
1271 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1272 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1273 assert!(!identity.records.is_empty());
1274 }
1275
1276 fn tagged_envelope() -> CoordinatorDataPlanEnvelope {
1284 let mut envelope = envelope();
1285 for record in &mut envelope.coordinator_relations.as_mut().unwrap().records {
1286 let (group, tag) = if record.sample_id.as_str() == "S002" {
1287 ("B", "dirty")
1288 } else {
1289 ("A", "clean")
1290 };
1291 record
1292 .metadata
1293 .insert("group".to_string(), serde_json::json!(group));
1294 record.tags = vec![tag.to_string()];
1295 }
1296 envelope.validate().unwrap();
1297 envelope
1298 }
1299
1300 #[test]
1301 fn branch_view_by_metadata_filters_relations_natively() {
1302 use crate::coordinator::{
1303 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1304 };
1305
1306 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1307 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1308
1309 let by_metadata = |group: &str| CoordinatorBranchView {
1310 view_id: format!("branch_view:group_{group}"),
1311 branch_id: format!("branch:group_{group}"),
1312 mode: CoordinatorBranchViewMode::ByMetadata,
1313 selector: CoordinatorBranchViewSelector {
1314 metadata: BTreeMap::from([("group".to_string(), serde_json::json!(group))]),
1315 ..Default::default()
1316 },
1317 allow_overlap: false,
1318 metadata: BTreeMap::new(),
1319 };
1320
1321 let view_a = DataView {
1323 branch_view: Some(by_metadata("A")),
1324 include_augmented: true,
1325 ..Default::default()
1326 };
1327 let record_a = arena.make_view(data.handle.handle, &view_a).unwrap();
1328 let identity_a = arena.view_identity(record_a.handle.handle).unwrap();
1329 assert!(
1330 identity_a
1331 .records
1332 .iter()
1333 .all(|record| record.sample_id.as_str() == "S001"),
1334 "by_metadata group=A must include only S001 relations"
1335 );
1336 assert!(
1337 identity_a
1338 .records
1339 .iter()
1340 .any(|record| record.sample_id.as_str() == "S001"),
1341 "by_metadata group=A must keep the matching S001 relations"
1342 );
1343
1344 let view_b = DataView {
1346 branch_view: Some(by_metadata("B")),
1347 include_augmented: true,
1348 ..Default::default()
1349 };
1350 let record_b = arena.make_view(data.handle.handle, &view_b).unwrap();
1351 let identity_b = arena.view_identity(record_b.handle.handle).unwrap();
1352 assert!(
1353 identity_b
1354 .records
1355 .iter()
1356 .all(|record| record.sample_id.as_str() == "S002"),
1357 "by_metadata group=B must exclude S001 relations"
1358 );
1359 }
1360
1361 #[test]
1362 fn branch_view_by_tag_filters_relations_natively() {
1363 use crate::coordinator::{
1364 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1365 };
1366
1367 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1368 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1369
1370 let by_tag = |tag: &str| CoordinatorBranchView {
1371 view_id: format!("branch_view:tag_{tag}"),
1372 branch_id: format!("branch:tag_{tag}"),
1373 mode: CoordinatorBranchViewMode::ByTag,
1374 selector: CoordinatorBranchViewSelector {
1375 tags: vec![tag.to_string()],
1376 ..Default::default()
1377 },
1378 allow_overlap: false,
1379 metadata: BTreeMap::new(),
1380 };
1381
1382 let view_clean = DataView {
1384 branch_view: Some(by_tag("clean")),
1385 include_augmented: true,
1386 ..Default::default()
1387 };
1388 let record_clean = arena.make_view(data.handle.handle, &view_clean).unwrap();
1389 let identity_clean = arena.view_identity(record_clean.handle.handle).unwrap();
1390 assert!(
1391 identity_clean
1392 .records
1393 .iter()
1394 .all(|record| record.sample_id.as_str() == "S001"),
1395 "by_tag clean must include only S001 relations"
1396 );
1397
1398 let view_dirty = DataView {
1400 branch_view: Some(by_tag("dirty")),
1401 include_augmented: true,
1402 ..Default::default()
1403 };
1404 let record_dirty = arena.make_view(data.handle.handle, &view_dirty).unwrap();
1405 let identity_dirty = arena.view_identity(record_dirty.handle.handle).unwrap();
1406 assert!(
1407 identity_dirty
1408 .records
1409 .iter()
1410 .all(|record| record.sample_id.as_str() == "S002"),
1411 "by_tag dirty must exclude S001 relations"
1412 );
1413 }
1414
1415 #[test]
1416 fn branch_view_by_filter_filters_closed_metadata_and_tag_predicates() {
1417 use crate::coordinator::{
1418 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1419 };
1420
1421 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1422 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1423 let view = DataView {
1424 branch_view: Some(CoordinatorBranchView {
1425 view_id: "branch_view:filter_A_clean".to_string(),
1426 branch_id: "branch:filter_A_clean".to_string(),
1427 mode: CoordinatorBranchViewMode::ByFilter,
1428 selector: CoordinatorBranchViewSelector {
1429 filter: Some(serde_json::json!({
1430 "metadata_equals": {"group": "A"},
1431 "tags_all": ["clean"]
1432 })),
1433 ..Default::default()
1434 },
1435 allow_overlap: false,
1436 metadata: BTreeMap::new(),
1437 }),
1438 include_augmented: true,
1439 ..Default::default()
1440 };
1441 let record = arena.make_view(data.handle.handle, &view).unwrap();
1442 let identity = arena.view_identity(record.handle.handle).unwrap();
1443 assert!(
1444 identity
1445 .records
1446 .iter()
1447 .all(|relation| relation.sample_id.as_str() == "S001"),
1448 "closed by_filter must retain exactly the matching relations"
1449 );
1450 }
1451
1452 #[test]
1453 fn branch_view_empty_partition_intersection_raises_a_clear_error() {
1454 use crate::coordinator::{
1459 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1460 };
1461
1462 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1463 let envelope = tagged_envelope();
1464 let data = arena.materialize(&envelope, &request()).unwrap();
1465
1466 let s001_samples: Vec<SampleId> = envelope
1469 .coordinator_relations
1470 .as_ref()
1471 .unwrap()
1472 .records
1473 .iter()
1474 .filter(|record| record.sample_id.as_str() == "S001")
1475 .map(|record| record.sample_id.clone())
1476 .collect::<BTreeSet<_>>()
1477 .into_iter()
1478 .collect();
1479 assert!(!s001_samples.is_empty(), "fixture must carry S001 samples");
1480
1481 let view = DataView {
1482 branch_view: Some(CoordinatorBranchView {
1483 view_id: "branch_view:group_B".to_string(),
1484 branch_id: "branch:group_B".to_string(),
1485 mode: CoordinatorBranchViewMode::ByMetadata,
1486 selector: CoordinatorBranchViewSelector {
1487 metadata: BTreeMap::from([("group".to_string(), serde_json::json!("B"))]),
1488 ..Default::default()
1489 },
1490 allow_overlap: false,
1491 metadata: BTreeMap::new(),
1492 }),
1493 sample_ids: Some(s001_samples),
1494 include_augmented: true,
1495 ..Default::default()
1496 };
1497 let error = arena
1498 .make_view(data.handle.handle, &view)
1499 .unwrap_err()
1500 .to_string();
1501 assert!(
1502 error.contains("selected no coordinator relations"),
1503 "empty branch ∩ fold must raise a clear error: {error}"
1504 );
1505 }
1506
1507 #[test]
1508 fn branch_view_by_filter_refuses_unknown_predicates_before_selection() {
1509 use crate::coordinator::{
1510 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1511 };
1512
1513 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1514 let data = arena.materialize(&envelope(), &request()).unwrap();
1515 let view = DataView {
1516 branch_view: Some(CoordinatorBranchView {
1517 view_id: "branch_view:by_filter".to_string(),
1518 branch_id: "branch:0".to_string(),
1519 mode: CoordinatorBranchViewMode::ByFilter,
1520 selector: CoordinatorBranchViewSelector {
1521 filter: Some(serde_json::json!({"op": "always"})),
1522 ..Default::default()
1523 },
1524 allow_overlap: false,
1525 metadata: BTreeMap::new(),
1526 }),
1527 include_augmented: true,
1528 ..Default::default()
1529 };
1530 let error = arena
1531 .make_view(data.handle.handle, &view)
1532 .expect_err("unknown by_filter predicates must fail closed");
1533 let message = format!("{error}");
1534 assert!(
1535 message.contains("native predicate"),
1536 "expected native-predicate validation error, got: {message}"
1537 );
1538 }
1539
1540 #[test]
1541 fn target_values_are_sample_level_and_dedup_repetitions() {
1542 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1543 let data = arena.materialize(&envelope(), &request()).unwrap();
1544 let view = DataView {
1545 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1546 include_augmented: false,
1547 ..Default::default()
1548 };
1549 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1550 let target_table = CoordinatorTargetTable {
1551 target_id: TargetId::new("y").unwrap(),
1552 values: vec![
1553 CoordinatorTargetValue {
1554 sample_id: SampleId::new("S001").unwrap(),
1555 value: json!(42.0),
1556 },
1557 CoordinatorTargetValue {
1558 sample_id: SampleId::new("S002").unwrap(),
1559 value: json!(7.0),
1560 },
1561 ],
1562 };
1563
1564 let target = arena
1565 .target_values(view_record.handle.handle, &target_table)
1566 .unwrap();
1567
1568 assert_eq!(target.target_id.as_str(), "y");
1569 assert_eq!(target.sample_ids, vec![SampleId::new("S001").unwrap()]);
1570 assert_eq!(target.values, vec![json!(42.0)]);
1571 }
1572
1573 #[test]
1574 fn multi_target_values_align_samples_and_emit_validity_masks() {
1575 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1576 let data = arena.materialize(&envelope(), &request()).unwrap();
1577 let view = DataView {
1578 sample_ids: Some(vec![
1579 SampleId::new("S002").unwrap(),
1580 SampleId::new("S001").unwrap(),
1581 ]),
1582 include_augmented: false,
1583 ..Default::default()
1584 };
1585 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1586 let y = CoordinatorTargetTable {
1587 target_id: TargetId::new("y").unwrap(),
1588 values: vec![
1589 CoordinatorTargetValue {
1590 sample_id: SampleId::new("S001").unwrap(),
1591 value: json!(42.0),
1592 },
1593 CoordinatorTargetValue {
1594 sample_id: SampleId::new("S002").unwrap(),
1595 value: json!(7.0),
1596 },
1597 ],
1598 };
1599 let protein = CoordinatorTargetTable {
1600 target_id: TargetId::new("protein").unwrap(),
1601 values: vec![CoordinatorTargetValue {
1602 sample_id: SampleId::new("S001").unwrap(),
1603 value: json!(12.5),
1604 }],
1605 };
1606
1607 let block = arena
1608 .multi_target_values(view_record.handle.handle, &[y, protein])
1609 .unwrap();
1610
1611 assert_eq!(
1612 block.target_ids,
1613 vec![
1614 TargetId::new("y").unwrap(),
1615 TargetId::new("protein").unwrap()
1616 ]
1617 );
1618 assert_eq!(
1619 block.sample_ids,
1620 vec![
1621 SampleId::new("S002").unwrap(),
1622 SampleId::new("S001").unwrap()
1623 ]
1624 );
1625 assert_eq!(block.values[0], vec![json!(7.0), json!(42.0)]);
1626 assert_eq!(block.validity_masks[0], vec![true, true]);
1627 assert_eq!(block.values[1], vec![json!(null), json!(12.5)]);
1628 assert_eq!(block.validity_masks[1], vec![false, true]);
1629 }
1630
1631 #[test]
1632 fn feature_values_are_observation_level_and_filter_columns() {
1633 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1634 let data = arena.materialize(&envelope(), &request()).unwrap();
1635 let view = DataView {
1636 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1637 columns: Some(vec!["f1".to_string()]),
1638 include_augmented: false,
1639 ..Default::default()
1640 };
1641 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1642 let feature_table = CoordinatorFeatureTable {
1643 feature_set_id: "x".to_string(),
1644 representation_id: RepresentationId::new("tabular_numeric").unwrap(),
1645 feature_names: vec!["f0".to_string(), "f1".to_string()],
1646 rows: vec![
1647 CoordinatorFeatureRow {
1648 observation_id: ObservationId::new("obs.S001.base").unwrap(),
1649 values: vec![json!(1.0), json!(10.0)],
1650 },
1651 CoordinatorFeatureRow {
1652 observation_id: ObservationId::new("obs.S001.rep1").unwrap(),
1653 values: vec![json!(2.0), json!(20.0)],
1654 },
1655 CoordinatorFeatureRow {
1656 observation_id: ObservationId::new("obs.S001.aug0").unwrap(),
1657 values: vec![json!(3.0), json!(30.0)],
1658 },
1659 CoordinatorFeatureRow {
1660 observation_id: ObservationId::new("obs.S002.base").unwrap(),
1661 values: vec![json!(4.0), json!(40.0)],
1662 },
1663 ],
1664 };
1665
1666 let features = arena
1667 .feature_values(view_record.handle.handle, &feature_table)
1668 .unwrap();
1669
1670 assert_eq!(features.feature_set_id, "x");
1671 assert_eq!(features.feature_names, vec!["f1".to_string()]);
1672 assert_eq!(features.representation_id.as_str(), "tabular_numeric");
1673 assert_eq!(
1674 features.observation_ids,
1675 vec![
1676 ObservationId::new("obs.S001.base").unwrap(),
1677 ObservationId::new("obs.S001.rep1").unwrap(),
1678 ]
1679 );
1680 assert_eq!(
1681 features.sample_ids,
1682 vec![
1683 SampleId::new("S001").unwrap(),
1684 SampleId::new("S001").unwrap()
1685 ]
1686 );
1687 assert_eq!(features.values, vec![vec![json!(10.0)], vec![json!(20.0)]]);
1688
1689 let mut wrong_representation = feature_table;
1690 wrong_representation.representation_id = RepresentationId::new("dense_signal").unwrap();
1691 assert!(arena
1692 .feature_values(view_record.handle.handle, &wrong_representation)
1693 .is_err());
1694 }
1695
1696 #[test]
1697 fn view_honors_requested_sample_order_for_identity_targets_and_features() {
1698 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1699 let data = arena.materialize(&envelope(), &request()).unwrap();
1700 let view = DataView {
1701 sample_ids: Some(vec![
1702 SampleId::new("S002").unwrap(),
1703 SampleId::new("S001").unwrap(),
1704 ]),
1705 include_augmented: false,
1706 ..Default::default()
1707 };
1708 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1709
1710 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1711 assert_eq!(
1712 identity
1713 .records
1714 .iter()
1715 .map(|relation| relation.observation_id.as_str())
1716 .collect::<Vec<_>>(),
1717 vec!["obs.S002.base", "obs.S001.base", "obs.S001.rep1"]
1718 );
1719
1720 let target_table = CoordinatorTargetTable {
1721 target_id: TargetId::new("y").unwrap(),
1722 values: vec![
1723 CoordinatorTargetValue {
1724 sample_id: SampleId::new("S001").unwrap(),
1725 value: json!(42.0),
1726 },
1727 CoordinatorTargetValue {
1728 sample_id: SampleId::new("S002").unwrap(),
1729 value: json!(7.0),
1730 },
1731 ],
1732 };
1733 let target = arena
1734 .target_values(view_record.handle.handle, &target_table)
1735 .unwrap();
1736 assert_eq!(
1737 target.sample_ids,
1738 vec![
1739 SampleId::new("S002").unwrap(),
1740 SampleId::new("S001").unwrap()
1741 ]
1742 );
1743 assert_eq!(target.values, vec![json!(7.0), json!(42.0)]);
1744
1745 let feature_table = CoordinatorFeatureTable {
1746 feature_set_id: "x".to_string(),
1747 representation_id: RepresentationId::new("tabular_numeric").unwrap(),
1748 feature_names: vec!["f0".to_string(), "f1".to_string()],
1749 rows: vec![
1750 CoordinatorFeatureRow {
1751 observation_id: ObservationId::new("obs.S001.base").unwrap(),
1752 values: vec![json!(1.0), json!(10.0)],
1753 },
1754 CoordinatorFeatureRow {
1755 observation_id: ObservationId::new("obs.S001.rep1").unwrap(),
1756 values: vec![json!(2.0), json!(20.0)],
1757 },
1758 CoordinatorFeatureRow {
1759 observation_id: ObservationId::new("obs.S002.base").unwrap(),
1760 values: vec![json!(4.0), json!(40.0)],
1761 },
1762 ],
1763 };
1764 let features = arena
1765 .feature_values(view_record.handle.handle, &feature_table)
1766 .unwrap();
1767 assert_eq!(
1768 features.observation_ids,
1769 vec![
1770 ObservationId::new("obs.S002.base").unwrap(),
1771 ObservationId::new("obs.S001.base").unwrap(),
1772 ObservationId::new("obs.S001.rep1").unwrap(),
1773 ]
1774 );
1775 assert_eq!(
1776 features.values,
1777 vec![
1778 vec![json!(4.0), json!(40.0)],
1779 vec![json!(1.0), json!(10.0)],
1780 vec![json!(2.0), json!(20.0)],
1781 ]
1782 );
1783 }
1784
1785 #[test]
1786 fn release_data_handle_releases_child_views() {
1787 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1788 let data = arena.materialize(&envelope(), &request()).unwrap();
1789 let view_record = arena
1790 .make_view(data.handle.handle, &DataView::default())
1791 .unwrap();
1792
1793 assert!(arena.release_handle(data.handle.handle));
1794 assert_eq!(arena.handle_record(data.handle.handle), None);
1795 assert_eq!(arena.view_record(view_record.handle.handle), None);
1796 let error = arena.view_identity(view_record.handle.handle).unwrap_err();
1797 assert_eq!(error.category(), "runtime");
1798 assert_eq!(error.code(), "unknown_handle");
1799 assert_eq!(error.error_code(), 0x0001_0001);
1800 assert!(!arena.release_handle(data.handle.handle));
1801 }
1802}