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