Skip to main content

dag_ml_data_core/
handle.rs

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    /// Target-major values: `values[target_idx][sample_idx]`.
158    pub values: Vec<Vec<serde_json::Value>>,
159    /// Target-major validity masks aligned with `values`.
160    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/// Typed sibling of [`CoordinatorFeatureBlock`]: the same projection with the
235/// cells flattened row-major into one `Vec<f64>` instead of per-cell JSON
236/// values. Produced by the `*_f64` projection path so a large numeric block
237/// never allocates `rows × cols` boxed [`serde_json::Value`]s. Masked cells
238/// are rejected by that path (it has no `Null` to project them into).
239#[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        // A handle is only "unknown" if it is absent from the data-handle
410        // registry. A live handle materialized without scoped relations is
411        // present in `records` but has no `data_relations` entry — that is a
412        // missing-relations contract error, not an unknown/stale handle.
413        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
810/// DAG owns fold membership. Explicit sample IDs are its authoritative
811/// selection and labels are provenance. Without IDs, named labels must match
812/// supplied relation metadata; only whole-handle full_train/predict views need
813/// no row predicate (PREDICT must already be bound to a predict handle).
814fn 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    // When both `view.source_ids` and `view.branch_view` constrain sources, the
858    // two are composed as intersection: a relation must pass both filters.
859    // `source_ids` is the materialization scope (sources loaded into the parent
860    // handle); the branch selector is the subset this branch cares about.
861    // Intersection prevents a branch from silently widening past what was
862    // actually materialized.
863    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            // A relation matches a `by_metadata` selector iff its own metadata
879            // contains every selector key with an equal value; `by_tag` iff its
880            // tags contain every selector tag.
881            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        // Mark the sole S002 observation excluded. Under a training policy
1249        // (`include_excluded=false`) it must be filtered out; under a
1250        // validation/predict policy (`include_excluded=true`) it must be kept.
1251        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        // `relation_fingerprint` is a replay key for the *source* relation
1258        // table, not for the derived coordinator_relations we just edited, so
1259        // it stays intact: the embedded coordinator relations are validated
1260        // structurally, and `excluded` does not participate in that fingerprint.
1261        envelope.validate().unwrap();
1262
1263        let arena = CoordinatorHandleArena::new("controller:data.provider").unwrap();
1264        let data = arena.materialize(&envelope, &request()).unwrap();
1265
1266        // Training view: excluded S002 row is dropped.
1267        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        // Validation/predict view: excluded S002 row is retained.
1284        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    /// Build the shared envelope, tagging S001 observations with
1405    /// `metadata={"group":"A"}` + `tags=["clean"]` and the S002 observation
1406    /// with `metadata={"group":"B"}` + `tags=["dirty"]`, so `by_metadata` /
1407    /// `by_tag` branch views have something to discriminate on. Editing the
1408    /// derived `coordinator_relations` does not invalidate the source
1409    /// `relation_fingerprint` (it is a replay key for the source table, not the
1410    /// derived view).
1411    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        // group=A INCLUDES the S001 relations and EXCLUDES S002.
1450        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        // group=B INCLUDES S002 and EXCLUDES S001.
1473        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        // tag=clean INCLUDES S001 and EXCLUDES S002.
1511        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        // tag=dirty INCLUDES S002 and EXCLUDES S001.
1527        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        // A branch_view scoped to group=B, intersected with a fold restriction
1583        // (sample_ids) that contains ONLY S001 samples, selects no relations —
1584        // the empty partition ∩ fold case. The arena must surface this explicitly
1585        // rather than returning an empty view (no silent mis-coverage).
1586        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        // The S001 sample id only (group=A); group=B (S002) is disjoint from this.
1595        // `DataView.sample_ids` rejects duplicates, so collect the distinct id.
1596        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}