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 selectors = resolved_view_labels(view, &parent)?;
356 let scoped = relations
357 .records
358 .into_iter()
359 .filter(|relation| {
360 selectors.iter().all(|(key, label)| {
361 relation
362 .metadata
363 .get(*key)
364 .and_then(serde_json::Value::as_str)
365 == Some(*label)
366 })
367 })
368 .collect::<Vec<_>>();
369 let filtered = filter_relations(&scoped, view)?;
370 let sample_count = unique_sample_count(&filtered);
371 let relation_record_count = filtered.len();
372 let handle = CoordinatorHandleRef {
373 handle: self.next_handle(),
374 kind: CoordinatorHandleKind::View,
375 owner_controller: self.owner_controller.clone(),
376 };
377 let record = CoordinatorDataViewRecord {
378 handle: handle.clone(),
379 parent_handle: parent.handle,
380 view: view.clone(),
381 sample_count,
382 relation_record_count,
383 };
384 self.view_records
385 .borrow_mut()
386 .insert(handle.handle, record.clone());
387 self.view_relations
388 .borrow_mut()
389 .insert(handle.handle, CoordinatorRelationSet { records: filtered });
390 Ok(record)
391 }
392
393 pub fn view_record(&self, handle: u64) -> Option<CoordinatorDataViewRecord> {
394 self.view_records.borrow().get(&handle).cloned()
395 }
396
397 pub fn view_identity(&self, handle: u64) -> Result<CoordinatorRelationSet> {
398 self.view_relations
399 .borrow()
400 .get(&handle)
401 .cloned()
402 .ok_or(DataError::UnknownHandle {
403 kind: "view",
404 handle,
405 })
406 }
407
408 pub fn data_identity(&self, handle: u64) -> Result<CoordinatorRelationSet> {
409 if !self.records.borrow().contains_key(&handle) {
414 return Err(DataError::UnknownHandle {
415 kind: "data",
416 handle,
417 });
418 }
419 self.data_relations
420 .borrow()
421 .get(&handle)
422 .cloned()
423 .ok_or_else(|| {
424 DataError::Validation(format!(
425 "data handle `{handle}` has no coordinator relations"
426 ))
427 })
428 }
429
430 pub fn release_handle(&self, handle: u64) -> bool {
431 if self.view_records.borrow_mut().remove(&handle).is_some() {
432 self.view_relations.borrow_mut().remove(&handle);
433 return true;
434 }
435 if let Some(record) = self.records.borrow_mut().remove(&handle) {
436 self.data_relations.borrow_mut().remove(&handle);
437 let child_views = self
438 .view_records
439 .borrow()
440 .iter()
441 .filter_map(|(view_handle, view_record)| {
442 (view_record.parent_handle == record.handle).then_some(*view_handle)
443 })
444 .collect::<Vec<_>>();
445 for view_handle in child_views {
446 self.view_records.borrow_mut().remove(&view_handle);
447 self.view_relations.borrow_mut().remove(&view_handle);
448 }
449 return true;
450 }
451 false
452 }
453
454 pub fn target_values(
455 &self,
456 view_handle: u64,
457 target_table: &CoordinatorTargetTable,
458 ) -> Result<CoordinatorTargetBlock> {
459 target_table.validate()?;
460 let relations = self.view_identity(view_handle)?;
461 let values_by_sample = target_table
462 .values
463 .iter()
464 .map(|value| (&value.sample_id, &value.value))
465 .collect::<BTreeMap<_, _>>();
466 let mut seen_samples = BTreeSet::new();
467 let mut sample_ids = Vec::new();
468 let mut values = Vec::new();
469 for relation in relations.records.iter().filter(|relation| {
470 relation
471 .target_id
472 .as_ref()
473 .map(|target_id| target_id == &target_table.target_id)
474 .unwrap_or(true)
475 }) {
476 if !seen_samples.insert(&relation.sample_id) {
477 continue;
478 }
479 let value = values_by_sample.get(&relation.sample_id).ok_or_else(|| {
480 DataError::Validation(format!(
481 "target table `{}` has no value for sample `{}`",
482 target_table.target_id, relation.sample_id
483 ))
484 })?;
485 sample_ids.push(relation.sample_id.clone());
486 values.push((*value).clone());
487 }
488 if sample_ids.is_empty() {
489 return Err(DataError::Validation(format!(
490 "view `{view_handle}` contains no samples for target `{}`",
491 target_table.target_id
492 )));
493 }
494 Ok(CoordinatorTargetBlock {
495 target_id: target_table.target_id.clone(),
496 sample_ids,
497 values,
498 })
499 }
500
501 pub fn multi_target_values(
502 &self,
503 view_handle: u64,
504 target_tables: &[CoordinatorTargetTable],
505 ) -> Result<CoordinatorMultiTargetBlock> {
506 if target_tables.is_empty() {
507 return Err(DataError::Validation(
508 "multi-target materialization requires at least one target table".to_string(),
509 ));
510 }
511 let mut seen_targets = BTreeSet::new();
512 for table in target_tables {
513 table.validate()?;
514 if !seen_targets.insert(table.target_id.clone()) {
515 return Err(DataError::Validation(format!(
516 "multi-target materialization contains duplicate target `{}`",
517 table.target_id
518 )));
519 }
520 }
521
522 let target_ids = target_tables
523 .iter()
524 .map(|table| table.target_id.clone())
525 .collect::<Vec<_>>();
526 let target_universe = target_ids.iter().collect::<BTreeSet<_>>();
527 let relations = self.view_identity(view_handle)?;
528 let mut seen_samples = BTreeSet::new();
529 let mut sample_ids = Vec::new();
530 for relation in relations.records.iter().filter(|relation| {
531 relation
532 .target_id
533 .as_ref()
534 .map(|target_id| target_universe.contains(target_id))
535 .unwrap_or(true)
536 }) {
537 if seen_samples.insert(relation.sample_id.clone()) {
538 sample_ids.push(relation.sample_id.clone());
539 }
540 }
541 if sample_ids.is_empty() {
542 return Err(DataError::Validation(format!(
543 "view `{view_handle}` contains no samples for requested targets"
544 )));
545 }
546
547 let mut values = Vec::with_capacity(target_tables.len());
548 let mut validity_masks = Vec::with_capacity(target_tables.len());
549 for table in target_tables {
550 let values_by_sample = table
551 .values
552 .iter()
553 .map(|value| (&value.sample_id, &value.value))
554 .collect::<BTreeMap<_, _>>();
555 let mut target_values = Vec::with_capacity(sample_ids.len());
556 let mut validity = Vec::with_capacity(sample_ids.len());
557 for sample_id in &sample_ids {
558 match values_by_sample.get(sample_id) {
559 Some(value) if !value.is_null() => {
560 target_values.push((*value).clone());
561 validity.push(true);
562 }
563 Some(_) | None => {
564 target_values.push(serde_json::Value::Null);
565 validity.push(false);
566 }
567 }
568 }
569 values.push(target_values);
570 validity_masks.push(validity);
571 }
572
573 Ok(CoordinatorMultiTargetBlock {
574 target_ids,
575 sample_ids,
576 values,
577 validity_masks,
578 })
579 }
580
581 pub fn feature_values(
582 &self,
583 view_handle: u64,
584 feature_table: &CoordinatorFeatureTable,
585 ) -> Result<CoordinatorFeatureBlock> {
586 feature_table.validate()?;
587 let view_record = self
588 .view_records
589 .borrow()
590 .get(&view_handle)
591 .cloned()
592 .ok_or(DataError::UnknownHandle {
593 kind: "view",
594 handle: view_handle,
595 })?;
596 let parent_record = self
597 .records
598 .borrow()
599 .get(&view_record.parent_handle.handle)
600 .cloned()
601 .ok_or(DataError::UnknownHandle {
602 kind: "data",
603 handle: view_record.parent_handle.handle,
604 })?;
605 if feature_table.representation_id != parent_record.output_representation {
606 return Err(DataError::Validation(format!(
607 "feature table `{}` representation `{}` does not match materialized output representation `{}`",
608 feature_table.feature_set_id,
609 feature_table.representation_id,
610 parent_record.output_representation
611 )));
612 }
613 let relations = self.view_identity(view_handle)?;
614 let selected_indices = selected_feature_indices(feature_table, &view_record.view)?;
615 let rows_by_observation = feature_table
616 .rows
617 .iter()
618 .map(|row| (&row.observation_id, row))
619 .collect::<BTreeMap<_, _>>();
620 let mut observation_ids = Vec::with_capacity(relations.records.len());
621 let mut sample_ids = Vec::with_capacity(relations.records.len());
622 let mut values = Vec::with_capacity(relations.records.len());
623 for relation in &relations.records {
624 let row = rows_by_observation
625 .get(&relation.observation_id)
626 .ok_or_else(|| {
627 DataError::Validation(format!(
628 "feature table `{}` has no row for observation `{}`",
629 feature_table.feature_set_id, relation.observation_id
630 ))
631 })?;
632 observation_ids.push(relation.observation_id.clone());
633 sample_ids.push(relation.sample_id.clone());
634 values.push(
635 selected_indices
636 .iter()
637 .map(|idx| row.values[*idx].clone())
638 .collect(),
639 );
640 }
641 Ok(CoordinatorFeatureBlock {
642 feature_set_id: feature_table.feature_set_id.clone(),
643 representation_id: feature_table.representation_id.clone(),
644 feature_names: selected_indices
645 .iter()
646 .map(|idx| feature_table.feature_names[*idx].clone())
647 .collect(),
648 observation_ids,
649 sample_ids,
650 values,
651 })
652 }
653
654 pub fn handle_record(&self, handle: u64) -> Option<CoordinatorDataHandleRecord> {
655 self.records.borrow().get(&handle).cloned()
656 }
657
658 pub fn handle_records(&self) -> Vec<CoordinatorDataHandleRecord> {
659 self.records.borrow().values().cloned().collect()
660 }
661
662 fn next_handle(&self) -> u64 {
663 let mut next = self.next_handle.borrow_mut();
664 let handle = *next;
665 *next += 1;
666 handle
667 }
668}
669
670fn validate_request_against_envelope(
671 envelope: &CoordinatorDataPlanEnvelope,
672 request: &CoordinatorDataMaterializationRequest,
673) -> Result<()> {
674 if request.schema_fingerprint != envelope.schema_fingerprint {
675 return Err(DataError::FingerprintMismatch {
676 kind: "schema",
677 expected: envelope.schema_fingerprint.clone(),
678 actual: request.schema_fingerprint.clone(),
679 });
680 }
681 if request.plan_fingerprint != envelope.plan_fingerprint {
682 return Err(DataError::FingerprintMismatch {
683 kind: "plan",
684 expected: envelope.plan_fingerprint.clone(),
685 actual: request.plan_fingerprint.clone(),
686 });
687 }
688 if request.relation_fingerprint != envelope.relation_fingerprint {
689 let none = || "<none>".to_string();
690 return Err(DataError::FingerprintMismatch {
691 kind: "relation",
692 expected: envelope.relation_fingerprint.clone().unwrap_or_else(none),
693 actual: request.relation_fingerprint.clone().unwrap_or_else(none),
694 });
695 }
696 if request.require_relations && envelope.coordinator_relations.is_none() {
697 return Err(DataError::Validation(format!(
698 "materialization request `{}` on `{}` requires coordinator relations",
699 request.input_name, request.node_id
700 )));
701 }
702 if request.output_representation != envelope.plan.output_representation {
703 return Err(DataError::Validation(format!(
704 "materialization request `{}` on `{}` output representation `{}` does not match plan output `{}`",
705 request.input_name,
706 request.node_id,
707 request.output_representation,
708 envelope.plan.output_representation
709 )));
710 }
711 if !request.source_ids.is_empty() {
712 let plan_sources = envelope
713 .plan
714 .steps
715 .iter()
716 .filter_map(|step| step.source_id.as_ref())
717 .collect::<BTreeSet<_>>();
718 for source_id in &request.source_ids {
719 if !plan_sources.contains(source_id) {
720 return Err(DataError::Validation(format!(
721 "materialization request `{}` on `{}` source `{}` is not present in data plan `{}`",
722 request.input_name, request.node_id, source_id, envelope.plan.id
723 )));
724 }
725 }
726 }
727 Ok(())
728}
729
730fn scoped_relations_for_materialization(
731 relations: &CoordinatorRelationSet,
732 request: &CoordinatorDataMaterializationRequest,
733) -> Result<CoordinatorRelationSet> {
734 relations.validate()?;
735 if request.source_ids.is_empty() {
736 return Ok(relations.clone());
737 }
738 let source_filter = request.source_ids.iter().collect::<BTreeSet<_>>();
739 let scoped = relations
740 .records
741 .iter()
742 .filter(|relation| {
743 relation
744 .source_id
745 .as_ref()
746 .map(|source_id| source_filter.contains(source_id))
747 .unwrap_or(false)
748 })
749 .cloned()
750 .collect::<Vec<_>>();
751 if scoped.is_empty() {
752 return Err(DataError::Validation(format!(
753 "materialization request `{}` on `{}` selected no coordinator relations for requested source ids",
754 request.input_name, request.node_id
755 )));
756 }
757 let scoped = CoordinatorRelationSet { records: scoped };
758 scoped.validate()?;
759 Ok(scoped)
760}
761
762fn validate_non_empty(label: &str, value: &str) -> Result<()> {
763 if value.trim().is_empty() {
764 return Err(DataError::Validation(format!("{label} must not be empty")));
765 }
766 Ok(())
767}
768
769fn validate_view(view: &DataView) -> Result<()> {
770 for (key, value) in [("partition", &view.partition), ("fold_id", &view.fold_id)] {
771 if let Some(value) = value {
772 validate_non_empty(key, value)?;
773 }
774 }
775 if let Some(samples) = &view.sample_ids {
776 let unique = samples.iter().collect::<BTreeSet<_>>();
777 if unique.len() != samples.len() {
778 return Err(DataError::Validation(
779 "data view contains duplicate sample ids".to_string(),
780 ));
781 }
782 }
783 if let Some(sources) = &view.source_ids {
784 let unique = sources.iter().collect::<BTreeSet<_>>();
785 if unique.len() != sources.len() {
786 return Err(DataError::Validation(
787 "data view contains duplicate source ids".to_string(),
788 ));
789 }
790 }
791 if let Some(columns) = &view.columns {
792 let unique = columns.iter().collect::<BTreeSet<_>>();
793 if unique.len() != columns.len() {
794 return Err(DataError::Validation(
795 "data view contains duplicate columns".to_string(),
796 ));
797 }
798 if columns.iter().any(|column| column.trim().is_empty()) {
799 return Err(DataError::Validation(
800 "data view contains an empty column".to_string(),
801 ));
802 }
803 }
804 if let Some(branch_view) = &view.branch_view {
805 branch_view.validate()?;
806 }
807 Ok(())
808}
809
810fn resolved_view_labels<'a>(
815 view: &'a DataView,
816 parent: &CoordinatorDataHandleRecord,
817) -> Result<Vec<(&'static str, &'a str)>> {
818 if view.sample_ids.is_some() {
819 return Ok(Vec::new());
820 }
821 let whole_handle = match view.partition.as_deref() {
822 Some("full_train") => true,
823 Some("predict") => {
824 if !parent.phase.eq_ignore_ascii_case("predict") {
825 return Err(DataError::Validation("partition=predict requires a PREDICT materialized handle or explicit sample_ids".into()));
826 }
827 true
828 }
829 _ => false,
830 };
831 let mut selectors = Vec::new();
832 if !whole_handle {
833 if let Some(partition) = &view.partition {
834 selectors.push(("partition", partition.as_str()));
835 }
836 }
837 if let Some(fold) = &view.fold_id {
838 if !whole_handle || parent.fold_id.as_ref() != Some(fold) {
839 selectors.push(("fold_id", fold.as_str()));
840 }
841 }
842 Ok(selectors)
843}
844
845fn filter_relations(
846 relations: &[CoordinatorRelation],
847 view: &DataView,
848) -> Result<Vec<CoordinatorRelation>> {
849 let sample_filter = view
850 .sample_ids
851 .as_ref()
852 .map(|sample_ids| sample_ids.iter().collect::<BTreeSet<_>>());
853 let source_filter = view
854 .source_ids
855 .as_ref()
856 .map(|source_ids| source_ids.iter().collect::<BTreeSet<_>>());
857 let mut branch_source_filter: Option<BTreeSet<&SourceId>> = None;
864 let mut branch_metadata_filter: Option<&BTreeMap<String, serde_json::Value>> = None;
865 let mut branch_tag_filter: Option<BTreeSet<&String>> = None;
866 let mut branch_native_filter = None;
867 if let Some(branch_view) = view.branch_view.as_ref() {
868 match branch_view.mode {
869 crate::coordinator::CoordinatorBranchViewMode::BySource => {
870 branch_source_filter = Some(
871 branch_view
872 .selector
873 .source_ids
874 .iter()
875 .collect::<BTreeSet<_>>(),
876 );
877 }
878 crate::coordinator::CoordinatorBranchViewMode::ByMetadata => {
882 branch_metadata_filter = Some(&branch_view.selector.metadata);
883 }
884 crate::coordinator::CoordinatorBranchViewMode::ByTag => {
885 branch_tag_filter = Some(branch_view.selector.tags.iter().collect::<BTreeSet<_>>());
886 }
887 crate::coordinator::CoordinatorBranchViewMode::Separation => {}
888 crate::coordinator::CoordinatorBranchViewMode::ByFilter => {
889 let filter = branch_view.selector.filter.as_ref().ok_or_else(|| {
890 DataError::Validation(format!(
891 "coordinator branch view `{}` mode=by_filter requires filter",
892 branch_view.view_id
893 ))
894 })?;
895 branch_native_filter = Some(parse_native_branch_view_filter(
896 filter,
897 &format!("coordinator branch view `{}`", branch_view.view_id),
898 )?);
899 }
900 }
901 }
902 let mut filtered = relations
903 .iter()
904 .enumerate()
905 .filter(|relation| {
906 let relation = relation.1;
907 sample_filter
908 .as_ref()
909 .map(|samples| samples.contains(&relation.sample_id))
910 .unwrap_or(true)
911 })
912 .filter(|relation| {
913 let relation = relation.1;
914 source_filter
915 .as_ref()
916 .map(|sources| {
917 relation
918 .source_id
919 .as_ref()
920 .map(|source_id| sources.contains(source_id))
921 .unwrap_or(false)
922 })
923 .unwrap_or(true)
924 })
925 .filter(|relation| {
926 let relation = relation.1;
927 branch_source_filter
928 .as_ref()
929 .map(|sources| {
930 relation
931 .source_id
932 .as_ref()
933 .map(|source_id| sources.contains(source_id))
934 .unwrap_or(false)
935 })
936 .unwrap_or(true)
937 })
938 .filter(|relation| {
939 let relation = relation.1;
940 branch_metadata_filter
941 .map(|selector| {
942 selector
943 .iter()
944 .all(|(key, value)| relation.metadata.get(key) == Some(value))
945 })
946 .unwrap_or(true)
947 })
948 .filter(|relation| {
949 let relation = relation.1;
950 branch_tag_filter
951 .as_ref()
952 .map(|tags| tags.iter().all(|tag| relation.tags.contains(tag)))
953 .unwrap_or(true)
954 })
955 .filter(|relation| {
956 let relation = relation.1;
957 branch_native_filter
958 .as_ref()
959 .map(|filter| {
960 filter
961 .metadata_equals
962 .iter()
963 .all(|(key, value)| relation.metadata.get(key) == Some(value))
964 && filter
965 .tags_all
966 .iter()
967 .all(|tag| relation.tags.contains(tag))
968 })
969 .unwrap_or(true)
970 })
971 .filter(|relation| view.include_augmented || !relation.1.is_augmented)
972 .filter(|relation| view.include_excluded || !relation.1.excluded)
973 .map(|(idx, relation)| (idx, relation.clone()))
974 .collect::<Vec<_>>();
975 if filtered.is_empty() {
976 return Err(DataError::Validation(
977 "data view selected no coordinator relations".to_string(),
978 ));
979 }
980 if let Some(sample_ids) = &view.sample_ids {
981 let sample_order = sample_ids
982 .iter()
983 .enumerate()
984 .map(|(idx, sample_id)| (sample_id, idx))
985 .collect::<BTreeMap<_, _>>();
986 filtered.sort_by_key(|(idx, relation)| {
987 (
988 sample_order
989 .get(&relation.sample_id)
990 .copied()
991 .unwrap_or(usize::MAX),
992 *idx,
993 )
994 });
995 }
996 Ok(filtered.into_iter().map(|(_, relation)| relation).collect())
997}
998
999fn unique_sample_count(relations: &[CoordinatorRelation]) -> usize {
1000 relations
1001 .iter()
1002 .map(|relation| &relation.sample_id)
1003 .collect::<BTreeSet<_>>()
1004 .len()
1005}
1006
1007fn selected_feature_indices(
1008 table: &CoordinatorFeatureTable,
1009 view: &DataView,
1010) -> Result<Vec<usize>> {
1011 let index_by_name = table
1012 .feature_names
1013 .iter()
1014 .enumerate()
1015 .map(|(idx, name)| (name, idx))
1016 .collect::<BTreeMap<_, _>>();
1017 let indices = if let Some(columns) = &view.columns {
1018 columns
1019 .iter()
1020 .map(|column| {
1021 index_by_name.get(column).copied().ok_or_else(|| {
1022 DataError::Validation(format!(
1023 "feature table `{}` has no feature column `{}`",
1024 table.feature_set_id, column
1025 ))
1026 })
1027 })
1028 .collect::<Result<Vec<_>>>()?
1029 } else {
1030 (0..table.feature_names.len()).collect()
1031 };
1032 if indices.is_empty() {
1033 return Err(DataError::Validation(format!(
1034 "feature table `{}` selected no feature columns",
1035 table.feature_set_id
1036 )));
1037 }
1038 Ok(indices)
1039}
1040
1041#[cfg(test)]
1042mod tests {
1043 use super::*;
1044 use serde_json::json;
1045
1046 fn envelope() -> CoordinatorDataPlanEnvelope {
1047 serde_json::from_str(include_str!(
1048 "../../../examples/fixtures/oof_campaign/coordinator_data_plan_envelope_nir.json"
1049 ))
1050 .unwrap()
1051 }
1052
1053 fn request() -> CoordinatorDataMaterializationRequest {
1054 serde_json::from_str(include_str!(
1055 "../../../examples/fixtures/oof_campaign/materialization_request_model_base_x.json"
1056 ))
1057 .unwrap()
1058 }
1059
1060 #[test]
1061 fn materializes_validated_coordinator_handle_record() {
1062 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1063 let record = arena.materialize(&envelope(), &request()).unwrap();
1064
1065 assert_eq!(record.handle.handle, 1);
1066 assert_eq!(record.handle.kind, CoordinatorHandleKind::Data);
1067 assert_eq!(record.input_name, "x");
1068 assert_eq!(record.plan_id, "nir-to-tabular");
1069 assert_eq!(record.sample_count, Some(2));
1070 assert_eq!(record.relation_record_count, Some(4));
1071 assert_eq!(arena.handle_record(1), Some(record));
1072 assert_eq!(arena.handle_records().len(), 1);
1073 }
1074
1075 #[test]
1076 fn view_labels_require_membership_but_preserve_host_resolved_views() {
1077 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1078 let record = arena.materialize(&envelope(), &request()).unwrap();
1079 for value in [
1080 json!({"partition":"does-not-exist"}),
1081 json!({"fold_id":"does-not-exist"}),
1082 json!({"partition":"fold_train"}),
1083 json!({"partition":"predict"}),
1084 ] {
1085 let view: DataView = serde_json::from_value(value).unwrap();
1086 assert!(arena.make_view(record.handle.handle, &view).is_err());
1087 }
1088 assert!(serde_json::from_value::<DataView>(json!({"sampl_ids":["S001"]})).is_err());
1089 let explicit: DataView = serde_json::from_value(
1090 json!({"sample_ids":["S001"],"partition":"fold_train","fold_id":"fold0"}),
1091 )
1092 .unwrap();
1093 assert_eq!(
1094 arena
1095 .make_view(record.handle.handle, &explicit)
1096 .unwrap()
1097 .sample_count,
1098 1
1099 );
1100 let all: DataView = serde_json::from_value(json!({"partition":"full_train"})).unwrap();
1101 assert_eq!(
1102 arena
1103 .make_view(record.handle.handle, &all)
1104 .unwrap()
1105 .sample_count,
1106 2
1107 );
1108 let mut predict_request = request();
1109 predict_request.phase = "PREDICT".into();
1110 let predict = arena.materialize(&envelope(), &predict_request).unwrap();
1111 let view: DataView = serde_json::from_value(json!({"partition":"predict"})).unwrap();
1112 assert_eq!(
1113 arena
1114 .make_view(predict.handle.handle, &view)
1115 .unwrap()
1116 .sample_count,
1117 2
1118 );
1119 }
1120
1121 #[test]
1122 fn view_resolves_partition_and_fold_from_supplied_relation_metadata() {
1123 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1124 let mut envelope = envelope();
1125 for relation in &mut envelope.coordinator_relations.as_mut().unwrap().records {
1126 relation.metadata.insert(
1127 "partition".into(),
1128 json!(if relation.sample_id.as_str() == "S001" {
1129 "train"
1130 } else {
1131 "validation"
1132 }),
1133 );
1134 relation.metadata.insert("fold_id".into(), json!("fold0"));
1135 }
1136 let record = arena.materialize(&envelope, &request()).unwrap();
1137 let view: DataView =
1138 serde_json::from_value(json!({"partition":"validation","fold_id":"fold0"})).unwrap();
1139 let selected = arena.make_view(record.handle.handle, &view).unwrap();
1140 assert_eq!(selected.sample_count, 1);
1141 assert!(arena
1142 .view_identity(selected.handle.handle)
1143 .unwrap()
1144 .records
1145 .iter()
1146 .all(|row| row.sample_id.as_str() == "S002"));
1147 }
1148
1149 #[test]
1150 fn materialization_refuses_fingerprint_mismatch() {
1151 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1152 let mut request = request();
1153 request.plan_fingerprint = "0".repeat(64);
1154
1155 assert!(arena.materialize(&envelope(), &request).is_err());
1156 }
1157
1158 #[test]
1159 fn materialization_scopes_relations_to_requested_sources() {
1160 let mut envelope = envelope();
1161 let chem = SourceId::new("chem").unwrap();
1162 envelope.plan.steps.push(crate::plan::DataPlanStep {
1163 kind: crate::plan::DataPlanStepKind::Materialize,
1164 source_id: Some(chem.clone()),
1165 adapter_id: None,
1166 input_representation: None,
1167 output_representation: Some(RepresentationId::new("tabular_numeric").unwrap()),
1168 fit_scope: crate::plan::FitScope::Stateless,
1169 requires_user_choice: false,
1170 metadata: BTreeMap::new(),
1171 });
1172 envelope.plan_fingerprint = crate::data_plan_fingerprint(&envelope.plan).unwrap();
1173 envelope.relation_fingerprint = None;
1174 envelope
1175 .coordinator_relations
1176 .as_mut()
1177 .unwrap()
1178 .records
1179 .push(CoordinatorRelation {
1180 observation_id: ObservationId::new("chem.S001").unwrap(),
1181 sample_id: SampleId::new("S001").unwrap(),
1182 target_id: Some(TargetId::new("y").unwrap()),
1183 group_id: None,
1184 origin_sample_id: None,
1185 source_id: Some(chem.clone()),
1186 is_augmented: false,
1187 excluded: false,
1188 metadata: BTreeMap::new(),
1189 tags: Vec::new(),
1190 });
1191 envelope.validate().unwrap();
1192
1193 let mut request = request();
1194 request.plan_fingerprint = envelope.plan_fingerprint.clone();
1195 request.relation_fingerprint = None;
1196 request.require_relations = false;
1197 request.source_ids = vec![SourceId::new("nir").unwrap()];
1198
1199 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1200 let data = arena.materialize(&envelope, &request).unwrap();
1201 let view = arena
1202 .make_view(data.handle.handle, &DataView::default())
1203 .unwrap();
1204 let identity = arena.view_identity(view.handle.handle).unwrap();
1205
1206 assert_eq!(data.relation_record_count, Some(4));
1207 assert_eq!(
1208 arena
1209 .data_identity(data.handle.handle)
1210 .unwrap()
1211 .records
1212 .len(),
1213 4
1214 );
1215 assert!(identity
1216 .records
1217 .iter()
1218 .all(|record| record.source_id.as_ref() == Some(&SourceId::new("nir").unwrap())));
1219 }
1220
1221 #[test]
1222 fn view_filters_augmented_rows_and_preserves_repetition_identity() {
1223 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1224 let data = arena.materialize(&envelope(), &request()).unwrap();
1225 let view = DataView {
1226 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1227 include_augmented: false,
1228 ..Default::default()
1229 };
1230
1231 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1232 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1233
1234 assert_eq!(view_record.handle.kind, CoordinatorHandleKind::View);
1235 assert_eq!(view_record.sample_count, 1);
1236 assert_eq!(view_record.relation_record_count, 2);
1237 assert_eq!(identity.records.len(), 2);
1238 assert_eq!(identity.records[0].observation_id.as_str(), "obs.S001.base");
1239 assert_eq!(identity.records[1].observation_id.as_str(), "obs.S001.rep1");
1240 assert_eq!(
1241 arena.view_record(view_record.handle.handle),
1242 Some(view_record)
1243 );
1244 }
1245
1246 #[test]
1247 fn view_drops_excluded_rows_for_training_and_keeps_them_otherwise() {
1248 let mut envelope = envelope();
1252 for record in &mut envelope.coordinator_relations.as_mut().unwrap().records {
1253 if record.observation_id.as_str() == "obs.S002.base" {
1254 record.excluded = true;
1255 }
1256 }
1257 envelope.validate().unwrap();
1262
1263 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1264 let data = arena.materialize(&envelope, &request()).unwrap();
1265
1266 let train_view = DataView {
1268 include_augmented: true,
1269 include_excluded: false,
1270 ..Default::default()
1271 };
1272 let train_record = arena.make_view(data.handle.handle, &train_view).unwrap();
1273 let train_identity = arena.view_identity(train_record.handle.handle).unwrap();
1274 assert!(
1275 train_identity
1276 .records
1277 .iter()
1278 .all(|record| record.sample_id.as_str() != "S002"),
1279 "excluded sample S002 must be absent from a training view"
1280 );
1281 assert_eq!(train_record.relation_record_count, 3);
1282
1283 let predict_view = DataView {
1285 include_augmented: true,
1286 include_excluded: true,
1287 ..Default::default()
1288 };
1289 let predict_record = arena.make_view(data.handle.handle, &predict_view).unwrap();
1290 let predict_identity = arena.view_identity(predict_record.handle.handle).unwrap();
1291 assert!(
1292 predict_identity
1293 .records
1294 .iter()
1295 .any(|record| record.sample_id.as_str() == "S002"),
1296 "excluded sample S002 must be present in a validation/predict view"
1297 );
1298 assert_eq!(predict_record.relation_record_count, 4);
1299 }
1300
1301 #[test]
1302 fn branch_view_by_source_filters_relations_to_branch_sources() {
1303 use crate::coordinator::{
1304 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1305 };
1306
1307 let mut envelope = envelope();
1308 let chem = SourceId::new("chem").unwrap();
1309 envelope.plan.steps.push(crate::plan::DataPlanStep {
1310 kind: crate::plan::DataPlanStepKind::Materialize,
1311 source_id: Some(chem.clone()),
1312 adapter_id: None,
1313 input_representation: None,
1314 output_representation: Some(RepresentationId::new("tabular_numeric").unwrap()),
1315 fit_scope: crate::plan::FitScope::Stateless,
1316 requires_user_choice: false,
1317 metadata: BTreeMap::new(),
1318 });
1319 envelope.plan_fingerprint = crate::data_plan_fingerprint(&envelope.plan).unwrap();
1320 envelope.relation_fingerprint = None;
1321 envelope
1322 .coordinator_relations
1323 .as_mut()
1324 .unwrap()
1325 .records
1326 .push(CoordinatorRelation {
1327 observation_id: ObservationId::new("chem.S001").unwrap(),
1328 sample_id: SampleId::new("S001").unwrap(),
1329 target_id: Some(TargetId::new("y").unwrap()),
1330 group_id: None,
1331 origin_sample_id: None,
1332 source_id: Some(chem.clone()),
1333 is_augmented: false,
1334 excluded: false,
1335 metadata: BTreeMap::new(),
1336 tags: Vec::new(),
1337 });
1338 envelope.validate().unwrap();
1339 let mut request = request();
1340 request.plan_fingerprint = envelope.plan_fingerprint.clone();
1341 request.relation_fingerprint = None;
1342 request.require_relations = false;
1343 request.source_ids = vec![SourceId::new("nir").unwrap(), chem.clone()];
1344
1345 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1346 let data = arena.materialize(&envelope, &request).unwrap();
1347 let view = DataView {
1348 branch_view: Some(CoordinatorBranchView {
1349 view_id: "branch_view:nir".to_string(),
1350 branch_id: "branch:nir_only".to_string(),
1351 mode: CoordinatorBranchViewMode::BySource,
1352 selector: CoordinatorBranchViewSelector {
1353 source_ids: vec![SourceId::new("nir").unwrap()],
1354 ..Default::default()
1355 },
1356 allow_overlap: false,
1357 metadata: BTreeMap::new(),
1358 }),
1359 include_augmented: true,
1360 ..Default::default()
1361 };
1362 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1363 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1364
1365 assert!(identity
1366 .records
1367 .iter()
1368 .all(|record| record.source_id.as_ref() == Some(&SourceId::new("nir").unwrap())));
1369 assert!(identity.records.iter().any(|record| record
1370 .observation_id
1371 .as_str()
1372 .starts_with("obs.S001")
1373 || record.observation_id.as_str().starts_with("obs.S002")));
1374 }
1375
1376 #[test]
1377 fn branch_view_separation_does_not_restrict_relations() {
1378 use crate::coordinator::{
1379 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1380 };
1381
1382 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1383 let data = arena.materialize(&envelope(), &request()).unwrap();
1384 let view = DataView {
1385 branch_view: Some(CoordinatorBranchView {
1386 view_id: "branch_view:separation".to_string(),
1387 branch_id: "branch:0".to_string(),
1388 mode: CoordinatorBranchViewMode::Separation,
1389 selector: CoordinatorBranchViewSelector {
1390 tags: vec!["clean".to_string()],
1391 ..Default::default()
1392 },
1393 allow_overlap: false,
1394 metadata: BTreeMap::new(),
1395 }),
1396 include_augmented: true,
1397 ..Default::default()
1398 };
1399 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1400 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1401 assert!(!identity.records.is_empty());
1402 }
1403
1404 fn tagged_envelope() -> CoordinatorDataPlanEnvelope {
1412 let mut envelope = envelope();
1413 for record in &mut envelope.coordinator_relations.as_mut().unwrap().records {
1414 let (group, tag) = if record.sample_id.as_str() == "S002" {
1415 ("B", "dirty")
1416 } else {
1417 ("A", "clean")
1418 };
1419 record
1420 .metadata
1421 .insert("group".to_string(), serde_json::json!(group));
1422 record.tags = vec![tag.to_string()];
1423 }
1424 envelope.validate().unwrap();
1425 envelope
1426 }
1427
1428 #[test]
1429 fn branch_view_by_metadata_filters_relations_natively() {
1430 use crate::coordinator::{
1431 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1432 };
1433
1434 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1435 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1436
1437 let by_metadata = |group: &str| CoordinatorBranchView {
1438 view_id: format!("branch_view:group_{group}"),
1439 branch_id: format!("branch:group_{group}"),
1440 mode: CoordinatorBranchViewMode::ByMetadata,
1441 selector: CoordinatorBranchViewSelector {
1442 metadata: BTreeMap::from([("group".to_string(), serde_json::json!(group))]),
1443 ..Default::default()
1444 },
1445 allow_overlap: false,
1446 metadata: BTreeMap::new(),
1447 };
1448
1449 let view_a = DataView {
1451 branch_view: Some(by_metadata("A")),
1452 include_augmented: true,
1453 ..Default::default()
1454 };
1455 let record_a = arena.make_view(data.handle.handle, &view_a).unwrap();
1456 let identity_a = arena.view_identity(record_a.handle.handle).unwrap();
1457 assert!(
1458 identity_a
1459 .records
1460 .iter()
1461 .all(|record| record.sample_id.as_str() == "S001"),
1462 "by_metadata group=A must include only S001 relations"
1463 );
1464 assert!(
1465 identity_a
1466 .records
1467 .iter()
1468 .any(|record| record.sample_id.as_str() == "S001"),
1469 "by_metadata group=A must keep the matching S001 relations"
1470 );
1471
1472 let view_b = DataView {
1474 branch_view: Some(by_metadata("B")),
1475 include_augmented: true,
1476 ..Default::default()
1477 };
1478 let record_b = arena.make_view(data.handle.handle, &view_b).unwrap();
1479 let identity_b = arena.view_identity(record_b.handle.handle).unwrap();
1480 assert!(
1481 identity_b
1482 .records
1483 .iter()
1484 .all(|record| record.sample_id.as_str() == "S002"),
1485 "by_metadata group=B must exclude S001 relations"
1486 );
1487 }
1488
1489 #[test]
1490 fn branch_view_by_tag_filters_relations_natively() {
1491 use crate::coordinator::{
1492 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1493 };
1494
1495 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1496 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1497
1498 let by_tag = |tag: &str| CoordinatorBranchView {
1499 view_id: format!("branch_view:tag_{tag}"),
1500 branch_id: format!("branch:tag_{tag}"),
1501 mode: CoordinatorBranchViewMode::ByTag,
1502 selector: CoordinatorBranchViewSelector {
1503 tags: vec![tag.to_string()],
1504 ..Default::default()
1505 },
1506 allow_overlap: false,
1507 metadata: BTreeMap::new(),
1508 };
1509
1510 let view_clean = DataView {
1512 branch_view: Some(by_tag("clean")),
1513 include_augmented: true,
1514 ..Default::default()
1515 };
1516 let record_clean = arena.make_view(data.handle.handle, &view_clean).unwrap();
1517 let identity_clean = arena.view_identity(record_clean.handle.handle).unwrap();
1518 assert!(
1519 identity_clean
1520 .records
1521 .iter()
1522 .all(|record| record.sample_id.as_str() == "S001"),
1523 "by_tag clean must include only S001 relations"
1524 );
1525
1526 let view_dirty = DataView {
1528 branch_view: Some(by_tag("dirty")),
1529 include_augmented: true,
1530 ..Default::default()
1531 };
1532 let record_dirty = arena.make_view(data.handle.handle, &view_dirty).unwrap();
1533 let identity_dirty = arena.view_identity(record_dirty.handle.handle).unwrap();
1534 assert!(
1535 identity_dirty
1536 .records
1537 .iter()
1538 .all(|record| record.sample_id.as_str() == "S002"),
1539 "by_tag dirty must exclude S001 relations"
1540 );
1541 }
1542
1543 #[test]
1544 fn branch_view_by_filter_filters_closed_metadata_and_tag_predicates() {
1545 use crate::coordinator::{
1546 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1547 };
1548
1549 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1550 let data = arena.materialize(&tagged_envelope(), &request()).unwrap();
1551 let view = DataView {
1552 branch_view: Some(CoordinatorBranchView {
1553 view_id: "branch_view:filter_A_clean".to_string(),
1554 branch_id: "branch:filter_A_clean".to_string(),
1555 mode: CoordinatorBranchViewMode::ByFilter,
1556 selector: CoordinatorBranchViewSelector {
1557 filter: Some(serde_json::json!({
1558 "metadata_equals": {"group": "A"},
1559 "tags_all": ["clean"]
1560 })),
1561 ..Default::default()
1562 },
1563 allow_overlap: false,
1564 metadata: BTreeMap::new(),
1565 }),
1566 include_augmented: true,
1567 ..Default::default()
1568 };
1569 let record = arena.make_view(data.handle.handle, &view).unwrap();
1570 let identity = arena.view_identity(record.handle.handle).unwrap();
1571 assert!(
1572 identity
1573 .records
1574 .iter()
1575 .all(|relation| relation.sample_id.as_str() == "S001"),
1576 "closed by_filter must retain exactly the matching relations"
1577 );
1578 }
1579
1580 #[test]
1581 fn branch_view_empty_partition_intersection_raises_a_clear_error() {
1582 use crate::coordinator::{
1587 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1588 };
1589
1590 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1591 let envelope = tagged_envelope();
1592 let data = arena.materialize(&envelope, &request()).unwrap();
1593
1594 let s001_samples: Vec<SampleId> = envelope
1597 .coordinator_relations
1598 .as_ref()
1599 .unwrap()
1600 .records
1601 .iter()
1602 .filter(|record| record.sample_id.as_str() == "S001")
1603 .map(|record| record.sample_id.clone())
1604 .collect::<BTreeSet<_>>()
1605 .into_iter()
1606 .collect();
1607 assert!(!s001_samples.is_empty(), "fixture must carry S001 samples");
1608
1609 let view = DataView {
1610 branch_view: Some(CoordinatorBranchView {
1611 view_id: "branch_view:group_B".to_string(),
1612 branch_id: "branch:group_B".to_string(),
1613 mode: CoordinatorBranchViewMode::ByMetadata,
1614 selector: CoordinatorBranchViewSelector {
1615 metadata: BTreeMap::from([("group".to_string(), serde_json::json!("B"))]),
1616 ..Default::default()
1617 },
1618 allow_overlap: false,
1619 metadata: BTreeMap::new(),
1620 }),
1621 sample_ids: Some(s001_samples),
1622 include_augmented: true,
1623 ..Default::default()
1624 };
1625 let error = arena
1626 .make_view(data.handle.handle, &view)
1627 .unwrap_err()
1628 .to_string();
1629 assert!(
1630 error.contains("selected no coordinator relations"),
1631 "empty branch ∩ fold must raise a clear error: {error}"
1632 );
1633 }
1634
1635 #[test]
1636 fn branch_view_by_filter_refuses_unknown_predicates_before_selection() {
1637 use crate::coordinator::{
1638 CoordinatorBranchView, CoordinatorBranchViewMode, CoordinatorBranchViewSelector,
1639 };
1640
1641 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1642 let data = arena.materialize(&envelope(), &request()).unwrap();
1643 let view = DataView {
1644 branch_view: Some(CoordinatorBranchView {
1645 view_id: "branch_view:by_filter".to_string(),
1646 branch_id: "branch:0".to_string(),
1647 mode: CoordinatorBranchViewMode::ByFilter,
1648 selector: CoordinatorBranchViewSelector {
1649 filter: Some(serde_json::json!({"op": "always"})),
1650 ..Default::default()
1651 },
1652 allow_overlap: false,
1653 metadata: BTreeMap::new(),
1654 }),
1655 include_augmented: true,
1656 ..Default::default()
1657 };
1658 let error = arena
1659 .make_view(data.handle.handle, &view)
1660 .expect_err("unknown by_filter predicates must fail closed");
1661 let message = format!("{error}");
1662 assert!(
1663 message.contains("native predicate"),
1664 "expected native-predicate validation error, got: {message}"
1665 );
1666 }
1667
1668 #[test]
1669 fn target_values_are_sample_level_and_dedup_repetitions() {
1670 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1671 let data = arena.materialize(&envelope(), &request()).unwrap();
1672 let view = DataView {
1673 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1674 include_augmented: false,
1675 ..Default::default()
1676 };
1677 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1678 let target_table = CoordinatorTargetTable {
1679 target_id: TargetId::new("y").unwrap(),
1680 values: vec![
1681 CoordinatorTargetValue {
1682 sample_id: SampleId::new("S001").unwrap(),
1683 value: json!(42.0),
1684 },
1685 CoordinatorTargetValue {
1686 sample_id: SampleId::new("S002").unwrap(),
1687 value: json!(7.0),
1688 },
1689 ],
1690 };
1691
1692 let target = arena
1693 .target_values(view_record.handle.handle, &target_table)
1694 .unwrap();
1695
1696 assert_eq!(target.target_id.as_str(), "y");
1697 assert_eq!(target.sample_ids, vec![SampleId::new("S001").unwrap()]);
1698 assert_eq!(target.values, vec![json!(42.0)]);
1699 }
1700
1701 #[test]
1702 fn multi_target_values_align_samples_and_emit_validity_masks() {
1703 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1704 let data = arena.materialize(&envelope(), &request()).unwrap();
1705 let view = DataView {
1706 sample_ids: Some(vec![
1707 SampleId::new("S002").unwrap(),
1708 SampleId::new("S001").unwrap(),
1709 ]),
1710 include_augmented: false,
1711 ..Default::default()
1712 };
1713 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1714 let y = CoordinatorTargetTable {
1715 target_id: TargetId::new("y").unwrap(),
1716 values: vec![
1717 CoordinatorTargetValue {
1718 sample_id: SampleId::new("S001").unwrap(),
1719 value: json!(42.0),
1720 },
1721 CoordinatorTargetValue {
1722 sample_id: SampleId::new("S002").unwrap(),
1723 value: json!(7.0),
1724 },
1725 ],
1726 };
1727 let protein = CoordinatorTargetTable {
1728 target_id: TargetId::new("protein").unwrap(),
1729 values: vec![CoordinatorTargetValue {
1730 sample_id: SampleId::new("S001").unwrap(),
1731 value: json!(12.5),
1732 }],
1733 };
1734
1735 let block = arena
1736 .multi_target_values(view_record.handle.handle, &[y, protein])
1737 .unwrap();
1738
1739 assert_eq!(
1740 block.target_ids,
1741 vec![
1742 TargetId::new("y").unwrap(),
1743 TargetId::new("protein").unwrap()
1744 ]
1745 );
1746 assert_eq!(
1747 block.sample_ids,
1748 vec![
1749 SampleId::new("S002").unwrap(),
1750 SampleId::new("S001").unwrap()
1751 ]
1752 );
1753 assert_eq!(block.values[0], vec![json!(7.0), json!(42.0)]);
1754 assert_eq!(block.validity_masks[0], vec![true, true]);
1755 assert_eq!(block.values[1], vec![json!(null), json!(12.5)]);
1756 assert_eq!(block.validity_masks[1], vec![false, true]);
1757 }
1758
1759 #[test]
1760 fn feature_values_are_observation_level_and_filter_columns() {
1761 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1762 let data = arena.materialize(&envelope(), &request()).unwrap();
1763 let view = DataView {
1764 sample_ids: Some(vec![SampleId::new("S001").unwrap()]),
1765 columns: Some(vec!["f1".to_string()]),
1766 include_augmented: false,
1767 ..Default::default()
1768 };
1769 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1770 let feature_table = CoordinatorFeatureTable {
1771 feature_set_id: "x".to_string(),
1772 representation_id: RepresentationId::new("tabular_numeric").unwrap(),
1773 feature_names: vec!["f0".to_string(), "f1".to_string()],
1774 rows: vec![
1775 CoordinatorFeatureRow {
1776 observation_id: ObservationId::new("obs.S001.base").unwrap(),
1777 values: vec![json!(1.0), json!(10.0)],
1778 },
1779 CoordinatorFeatureRow {
1780 observation_id: ObservationId::new("obs.S001.rep1").unwrap(),
1781 values: vec![json!(2.0), json!(20.0)],
1782 },
1783 CoordinatorFeatureRow {
1784 observation_id: ObservationId::new("obs.S001.aug0").unwrap(),
1785 values: vec![json!(3.0), json!(30.0)],
1786 },
1787 CoordinatorFeatureRow {
1788 observation_id: ObservationId::new("obs.S002.base").unwrap(),
1789 values: vec![json!(4.0), json!(40.0)],
1790 },
1791 ],
1792 };
1793
1794 let features = arena
1795 .feature_values(view_record.handle.handle, &feature_table)
1796 .unwrap();
1797
1798 assert_eq!(features.feature_set_id, "x");
1799 assert_eq!(features.feature_names, vec!["f1".to_string()]);
1800 assert_eq!(features.representation_id.as_str(), "tabular_numeric");
1801 assert_eq!(
1802 features.observation_ids,
1803 vec![
1804 ObservationId::new("obs.S001.base").unwrap(),
1805 ObservationId::new("obs.S001.rep1").unwrap(),
1806 ]
1807 );
1808 assert_eq!(
1809 features.sample_ids,
1810 vec![
1811 SampleId::new("S001").unwrap(),
1812 SampleId::new("S001").unwrap()
1813 ]
1814 );
1815 assert_eq!(features.values, vec![vec![json!(10.0)], vec![json!(20.0)]]);
1816
1817 let mut wrong_representation = feature_table;
1818 wrong_representation.representation_id = RepresentationId::new("dense_signal").unwrap();
1819 assert!(arena
1820 .feature_values(view_record.handle.handle, &wrong_representation)
1821 .is_err());
1822 }
1823
1824 #[test]
1825 fn view_honors_requested_sample_order_for_identity_targets_and_features() {
1826 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1827 let data = arena.materialize(&envelope(), &request()).unwrap();
1828 let view = DataView {
1829 sample_ids: Some(vec![
1830 SampleId::new("S002").unwrap(),
1831 SampleId::new("S001").unwrap(),
1832 ]),
1833 include_augmented: false,
1834 ..Default::default()
1835 };
1836 let view_record = arena.make_view(data.handle.handle, &view).unwrap();
1837
1838 let identity = arena.view_identity(view_record.handle.handle).unwrap();
1839 assert_eq!(
1840 identity
1841 .records
1842 .iter()
1843 .map(|relation| relation.observation_id.as_str())
1844 .collect::<Vec<_>>(),
1845 vec!["obs.S002.base", "obs.S001.base", "obs.S001.rep1"]
1846 );
1847
1848 let target_table = CoordinatorTargetTable {
1849 target_id: TargetId::new("y").unwrap(),
1850 values: vec![
1851 CoordinatorTargetValue {
1852 sample_id: SampleId::new("S001").unwrap(),
1853 value: json!(42.0),
1854 },
1855 CoordinatorTargetValue {
1856 sample_id: SampleId::new("S002").unwrap(),
1857 value: json!(7.0),
1858 },
1859 ],
1860 };
1861 let target = arena
1862 .target_values(view_record.handle.handle, &target_table)
1863 .unwrap();
1864 assert_eq!(
1865 target.sample_ids,
1866 vec![
1867 SampleId::new("S002").unwrap(),
1868 SampleId::new("S001").unwrap()
1869 ]
1870 );
1871 assert_eq!(target.values, vec![json!(7.0), json!(42.0)]);
1872
1873 let feature_table = CoordinatorFeatureTable {
1874 feature_set_id: "x".to_string(),
1875 representation_id: RepresentationId::new("tabular_numeric").unwrap(),
1876 feature_names: vec!["f0".to_string(), "f1".to_string()],
1877 rows: vec![
1878 CoordinatorFeatureRow {
1879 observation_id: ObservationId::new("obs.S001.base").unwrap(),
1880 values: vec![json!(1.0), json!(10.0)],
1881 },
1882 CoordinatorFeatureRow {
1883 observation_id: ObservationId::new("obs.S001.rep1").unwrap(),
1884 values: vec![json!(2.0), json!(20.0)],
1885 },
1886 CoordinatorFeatureRow {
1887 observation_id: ObservationId::new("obs.S002.base").unwrap(),
1888 values: vec![json!(4.0), json!(40.0)],
1889 },
1890 ],
1891 };
1892 let features = arena
1893 .feature_values(view_record.handle.handle, &feature_table)
1894 .unwrap();
1895 assert_eq!(
1896 features.observation_ids,
1897 vec![
1898 ObservationId::new("obs.S002.base").unwrap(),
1899 ObservationId::new("obs.S001.base").unwrap(),
1900 ObservationId::new("obs.S001.rep1").unwrap(),
1901 ]
1902 );
1903 assert_eq!(
1904 features.values,
1905 vec![
1906 vec![json!(4.0), json!(40.0)],
1907 vec![json!(1.0), json!(10.0)],
1908 vec![json!(2.0), json!(20.0)],
1909 ]
1910 );
1911 }
1912
1913 #[test]
1914 fn release_data_handle_releases_child_views() {
1915 let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1916 let data = arena.materialize(&envelope(), &request()).unwrap();
1917 let view_record = arena
1918 .make_view(data.handle.handle, &DataView::default())
1919 .unwrap();
1920
1921 assert!(arena.release_handle(data.handle.handle));
1922 assert_eq!(arena.handle_record(data.handle.handle), None);
1923 assert_eq!(arena.view_record(view_record.handle.handle), None);
1924 let error = arena.view_identity(view_record.handle.handle).unwrap_err();
1925 assert_eq!(error.category(), "runtime");
1926 assert_eq!(error.code(), "unknown_handle");
1927 assert_eq!(error.error_code(), 0x0001_0001);
1928 assert!(!arena.release_handle(data.handle.handle));
1929 }
1930}