1use std::collections::{BTreeMap, BTreeSet};
2
3use serde::{de::DeserializeOwned, Deserialize, Serialize};
4use serde_json::{json, Value};
5use sha2::{Digest, Sha256};
6
7use crate::bundle::ExecutionBundle;
8use crate::data::ExternalDataPlanEnvelope;
9use crate::error::{DagMlError, Result};
10use crate::ids::{ArtifactId, LineageId};
11use crate::plan::ExecutionPlan;
12use crate::runtime::{
13 FileArtifactManifest, FilePredictionCacheManifest, LineageRecord, FILE_ARTIFACT_MANIFEST_FILE,
14 FILE_PREDICTION_CACHE_MANIFEST_FILE,
15};
16
17pub const RESEARCH_PROVENANCE_SCHEMA_VERSION: u32 = 1;
18pub const EXECUTION_PLAN_FILE: &str = "execution_plan.json";
19pub const EXECUTION_BUNDLE_FILE: &str = "execution_bundle.json";
20pub const LINEAGE_RECORDS_FILE: &str = "lineage_records.json";
21pub const PROV_JSONLD_FILE: &str = "lineage.prov.jsonld";
22pub const RO_CRATE_METADATA_FILE: &str = "ro-crate-metadata.json";
23pub const OPENLINEAGE_RUN_EVENT_SCHEMA_URL: &str =
24 "https://openlineage.io/spec/1-0-0/OpenLineage.json#/definitions/RunEvent";
25pub const DAGML_OPENLINEAGE_FACET_SCHEMA_URL: &str =
26 "https://github.com/GBeurier/dag-ml/schemas/openlineage_dagml_facets.v1.schema.json";
27
28#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
29pub struct ResearchProvenanceExport {
30 pub schema_version: u32,
31 pub prov_jsonld: Value,
32 pub ro_crate_metadata: Value,
33}
34
35#[derive(Clone, Debug, Eq, PartialEq)]
36pub struct ResearchProvenancePackage {
37 pub schema_version: u32,
38 pub files: BTreeMap<String, ResearchProvenancePackageFile>,
39}
40
41#[derive(Clone, Debug, Eq, PartialEq)]
42pub struct ResearchProvenancePackageFile {
43 pub path: String,
44 pub sha256: String,
45 pub size_bytes: usize,
46 pub bytes: Vec<u8>,
47}
48
49#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
50pub struct ResearchProvenancePackageValidation {
51 pub schema_version: u32,
52 pub plan_id: String,
53 pub bundle_id: String,
54 pub file_count: usize,
55 pub checksummed_file_count: usize,
56 pub lineage_record_count: usize,
57 pub data_envelope_count: usize,
58 pub has_prediction_cache_manifest: bool,
59 pub has_artifact_manifest: bool,
60}
61
62#[derive(Clone, Debug, Eq, PartialEq)]
63pub struct OpenLineageRunEventOptions {
64 pub namespace: String,
65 pub event_time: String,
66}
67
68pub fn build_research_provenance_export(
69 plan: &ExecutionPlan,
70 bundle: &ExecutionBundle,
71 lineage: &[LineageRecord],
72 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
73 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
74 artifact_manifest: Option<&FileArtifactManifest>,
75) -> Result<ResearchProvenanceExport> {
76 validate_provenance_inputs(
77 plan,
78 bundle,
79 lineage,
80 data_envelopes,
81 prediction_cache_manifest,
82 artifact_manifest,
83 )?;
84
85 Ok(ResearchProvenanceExport {
86 schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
87 prov_jsonld: build_prov_jsonld(
88 plan,
89 bundle,
90 lineage,
91 data_envelopes,
92 prediction_cache_manifest,
93 artifact_manifest,
94 )?,
95 ro_crate_metadata: build_ro_crate_metadata(
96 plan,
97 bundle,
98 data_envelopes,
99 prediction_cache_manifest,
100 artifact_manifest,
101 )?,
102 })
103}
104
105pub fn build_research_provenance_package(
106 plan: &ExecutionPlan,
107 bundle: &ExecutionBundle,
108 lineage: &[LineageRecord],
109 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
110 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
111 artifact_manifest: Option<&FileArtifactManifest>,
112) -> Result<ResearchProvenancePackage> {
113 let export = build_research_provenance_export(
114 plan,
115 bundle,
116 lineage,
117 data_envelopes,
118 prediction_cache_manifest,
119 artifact_manifest,
120 )?;
121 let mut files = BTreeMap::new();
122 add_json_package_file(&mut files, EXECUTION_PLAN_FILE, plan, "execution plan")?;
123 add_json_package_file(
124 &mut files,
125 EXECUTION_BUNDLE_FILE,
126 bundle,
127 "execution bundle",
128 )?;
129 add_json_package_file(
130 &mut files,
131 LINEAGE_RECORDS_FILE,
132 &lineage,
133 "lineage records",
134 )?;
135 add_json_package_file(
136 &mut files,
137 PROV_JSONLD_FILE,
138 &export.prov_jsonld,
139 "PROV JSON-LD",
140 )?;
141 if let Some(manifest) = prediction_cache_manifest {
142 add_json_package_file(
143 &mut files,
144 FILE_PREDICTION_CACHE_MANIFEST_FILE,
145 manifest,
146 "prediction cache manifest",
147 )?;
148 }
149 if let Some(manifest) = artifact_manifest {
150 add_json_package_file(
151 &mut files,
152 FILE_ARTIFACT_MANIFEST_FILE,
153 manifest,
154 "artifact manifest",
155 )?;
156 }
157 for (key, envelope) in data_envelopes {
158 add_json_package_file(
159 &mut files,
160 &data_envelope_file_path(key)?,
161 envelope,
162 "data envelope",
163 )?;
164 }
165
166 let mut ro_crate_metadata = export.ro_crate_metadata;
167 annotate_ro_crate_package_files(&mut ro_crate_metadata, &files)?;
168 add_json_package_file(
169 &mut files,
170 RO_CRATE_METADATA_FILE,
171 &ro_crate_metadata,
172 "RO-Crate metadata",
173 )?;
174
175 Ok(ResearchProvenancePackage {
176 schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
177 files,
178 })
179}
180
181pub fn validate_research_provenance_package_files(
182 files: &BTreeMap<String, Vec<u8>>,
183) -> Result<ResearchProvenancePackageValidation> {
184 if files.is_empty() {
185 return Err(DagMlError::RuntimeValidation(
186 "research provenance package has no files".to_string(),
187 ));
188 }
189 for path in files.keys() {
190 validate_package_path(path)?;
191 }
192 require_package_file(files, EXECUTION_PLAN_FILE)?;
193 require_package_file(files, EXECUTION_BUNDLE_FILE)?;
194 require_package_file(files, LINEAGE_RECORDS_FILE)?;
195 require_package_file(files, PROV_JSONLD_FILE)?;
196 let ro_crate_metadata: Value = parse_package_json(
197 require_package_file(files, RO_CRATE_METADATA_FILE)?,
198 RO_CRATE_METADATA_FILE,
199 )?;
200
201 let checksummed_file_count = validate_ro_crate_package_checksums(&ro_crate_metadata, files)?;
202 validate_prov_jsonld_root(parse_package_json(
203 require_package_file(files, PROV_JSONLD_FILE)?,
204 PROV_JSONLD_FILE,
205 )?)?;
206
207 let plan: ExecutionPlan = parse_package_json(
208 require_package_file(files, EXECUTION_PLAN_FILE)?,
209 EXECUTION_PLAN_FILE,
210 )?;
211 let bundle: ExecutionBundle = parse_package_json(
212 require_package_file(files, EXECUTION_BUNDLE_FILE)?,
213 EXECUTION_BUNDLE_FILE,
214 )?;
215 let lineage: Vec<LineageRecord> = parse_package_json(
216 require_package_file(files, LINEAGE_RECORDS_FILE)?,
217 LINEAGE_RECORDS_FILE,
218 )?;
219 let data_envelopes = parse_package_data_envelopes(files)?;
220 let prediction_cache_manifest: Option<FilePredictionCacheManifest> = files
221 .get(FILE_PREDICTION_CACHE_MANIFEST_FILE)
222 .map(|bytes| parse_package_json(bytes, FILE_PREDICTION_CACHE_MANIFEST_FILE))
223 .transpose()?;
224 let artifact_manifest: Option<FileArtifactManifest> = files
225 .get(FILE_ARTIFACT_MANIFEST_FILE)
226 .map(|bytes| parse_package_json(bytes, FILE_ARTIFACT_MANIFEST_FILE))
227 .transpose()?;
228
229 validate_provenance_inputs(
230 &plan,
231 &bundle,
232 &lineage,
233 &data_envelopes,
234 prediction_cache_manifest.as_ref(),
235 artifact_manifest.as_ref(),
236 )?;
237
238 Ok(ResearchProvenancePackageValidation {
239 schema_version: RESEARCH_PROVENANCE_SCHEMA_VERSION,
240 plan_id: plan.id.to_string(),
241 bundle_id: bundle.bundle_id.to_string(),
242 file_count: files.len(),
243 checksummed_file_count,
244 lineage_record_count: lineage.len(),
245 data_envelope_count: data_envelopes.len(),
246 has_prediction_cache_manifest: prediction_cache_manifest.is_some(),
247 has_artifact_manifest: artifact_manifest.is_some(),
248 })
249}
250
251pub fn build_openlineage_run_event_from_package_files(
252 files: &BTreeMap<String, Vec<u8>>,
253 namespace: &str,
254 event_time: &str,
255) -> Result<Value> {
256 validate_research_provenance_package_files(files)?;
257 let plan: ExecutionPlan = parse_package_json(
258 require_package_file(files, EXECUTION_PLAN_FILE)?,
259 EXECUTION_PLAN_FILE,
260 )?;
261 let bundle: ExecutionBundle = parse_package_json(
262 require_package_file(files, EXECUTION_BUNDLE_FILE)?,
263 EXECUTION_BUNDLE_FILE,
264 )?;
265 let lineage: Vec<LineageRecord> = parse_package_json(
266 require_package_file(files, LINEAGE_RECORDS_FILE)?,
267 LINEAGE_RECORDS_FILE,
268 )?;
269 let data_envelopes = parse_package_data_envelopes(files)?;
270 let prediction_cache_manifest: Option<FilePredictionCacheManifest> = files
271 .get(FILE_PREDICTION_CACHE_MANIFEST_FILE)
272 .map(|bytes| parse_package_json(bytes, FILE_PREDICTION_CACHE_MANIFEST_FILE))
273 .transpose()?;
274 let artifact_manifest: Option<FileArtifactManifest> = files
275 .get(FILE_ARTIFACT_MANIFEST_FILE)
276 .map(|bytes| parse_package_json(bytes, FILE_ARTIFACT_MANIFEST_FILE))
277 .transpose()?;
278 let options = OpenLineageRunEventOptions {
279 namespace: namespace.to_string(),
280 event_time: event_time.to_string(),
281 };
282
283 build_openlineage_run_event(
284 &plan,
285 &bundle,
286 &lineage,
287 &data_envelopes,
288 prediction_cache_manifest.as_ref(),
289 artifact_manifest.as_ref(),
290 &options,
291 )
292}
293
294pub fn build_openlineage_run_event(
295 plan: &ExecutionPlan,
296 bundle: &ExecutionBundle,
297 lineage: &[LineageRecord],
298 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
299 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
300 artifact_manifest: Option<&FileArtifactManifest>,
301 options: &OpenLineageRunEventOptions,
302) -> Result<Value> {
303 validate_provenance_inputs(
304 plan,
305 bundle,
306 lineage,
307 data_envelopes,
308 prediction_cache_manifest,
309 artifact_manifest,
310 )?;
311 validate_openlineage_namespace(options.namespace.as_str())?;
312 validate_openlineage_event_time(options.event_time.as_str())?;
313
314 Ok(json!({
315 "eventType": "COMPLETE",
316 "eventTime": options.event_time.as_str(),
317 "run": {
318 "runId": openlineage_run_id(plan, bundle),
319 "facets": {
320 "dagml_reproducibility": dagml_openlineage_reproducibility_run_facet(plan, bundle),
321 "dagml_oof_safety": dagml_openlineage_oof_safety_run_facet(bundle, lineage),
322 },
323 },
324 "job": {
325 "namespace": options.namespace.as_str(),
326 "name": format!("{}::{}", plan.id, bundle.bundle_id),
327 "facets": {
328 "dagml_plan": dagml_openlineage_plan_job_facet(plan, bundle),
329 },
330 },
331 "inputs": openlineage_input_datasets(bundle, data_envelopes),
332 "outputs": openlineage_output_datasets(bundle, prediction_cache_manifest, artifact_manifest),
333 "producer": "https://github.com/GBeurier/dag-ml",
334 "schemaURL": OPENLINEAGE_RUN_EVENT_SCHEMA_URL,
335 }))
336}
337
338fn validate_provenance_inputs(
339 plan: &ExecutionPlan,
340 bundle: &ExecutionBundle,
341 lineage: &[LineageRecord],
342 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
343 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
344 artifact_manifest: Option<&FileArtifactManifest>,
345) -> Result<()> {
346 plan.validate()?;
347 bundle.validate_against_plan(plan)?;
348 if !data_envelopes.is_empty() {
349 bundle.validate_replay_envelopes(data_envelopes)?;
350 }
351 if let Some(manifest) = prediction_cache_manifest {
352 manifest.validate_against_bundle(bundle)?;
353 }
354 if let Some(manifest) = artifact_manifest {
355 manifest.validate_against_bundle(bundle)?;
356 }
357
358 let mut lineage_ids = BTreeSet::<&LineageId>::new();
359 for record in lineage {
360 record.validate()?;
361 if !plan.node_plans.contains_key(&record.node_id) {
362 return Err(DagMlError::RuntimeValidation(format!(
363 "provenance lineage `{}` references unknown node `{}`",
364 record.record_id, record.node_id
365 )));
366 }
367 if !plan
368 .controller_manifests
369 .contains_key(&record.controller_id)
370 {
371 return Err(DagMlError::RuntimeValidation(format!(
372 "provenance lineage `{}` references unknown controller `{}`",
373 record.record_id, record.controller_id
374 )));
375 }
376 if !lineage_ids.insert(&record.record_id) {
377 return Err(DagMlError::RuntimeValidation(format!(
378 "duplicate provenance lineage record `{}`",
379 record.record_id
380 )));
381 }
382 }
383 for record in lineage {
384 for input_id in &record.input_lineage {
385 if !lineage_ids.contains(input_id) {
386 return Err(DagMlError::RuntimeValidation(format!(
387 "provenance lineage `{}` references missing input lineage `{}`",
388 record.record_id, input_id
389 )));
390 }
391 }
392 }
393 Ok(())
394}
395
396fn build_prov_jsonld(
397 plan: &ExecutionPlan,
398 bundle: &ExecutionBundle,
399 lineage: &[LineageRecord],
400 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
401 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
402 artifact_manifest: Option<&FileArtifactManifest>,
403) -> Result<Value> {
404 let plan_entity_id = format!("dagml:execution-plan:{}", plan.id);
405 let bundle_entity_id = format!("dagml:execution-bundle:{}", bundle.bundle_id);
406 let packaging_activity_id = format!("dagml:activity:package-bundle:{}", bundle.bundle_id);
407 let coordinator_agent_id = "dagml:agent:dag-ml".to_string();
408
409 let mut entity = BTreeMap::<String, Value>::new();
410 entity.insert(
411 plan_entity_id.clone(),
412 json!({
413 "prov:type": ["prov:Entity", "dagml:ExecutionPlan"],
414 "dagml:plan_id": plan.id,
415 "dagml:graph_fingerprint": plan.graph_fingerprint,
416 "dagml:campaign_fingerprint": plan.campaign_fingerprint,
417 "dagml:controller_fingerprint": plan.controller_fingerprint,
418 "dagml:variant_count": plan.variants.len(),
419 "dagml:has_fold_set": plan.fold_set.is_some(),
420 }),
421 );
422 entity.insert(
423 bundle_entity_id.clone(),
424 json!({
425 "prov:type": ["prov:Entity", "dagml:ExecutionBundle"],
426 "dagml:bundle_id": bundle.bundle_id,
427 "dagml:schema_version": bundle.schema_version,
428 "dagml:plan_id": bundle.plan_id,
429 "dagml:selected_variant_id": bundle.selected_variant_id,
430 "dagml:graph_fingerprint": bundle.graph_fingerprint,
431 "dagml:campaign_fingerprint": bundle.campaign_fingerprint,
432 "dagml:controller_fingerprint": bundle.controller_fingerprint,
433 "dagml:unsafe_flags": bundle.unsafe_flags,
434 "dagml:selection_count": bundle.selections.len(),
435 }),
436 );
437
438 for requirement in &bundle.data_requirements {
439 let key = requirement.key();
440 entity.insert(
441 data_requirement_entity_id(&key),
442 json!({
443 "prov:type": ["prov:Entity", "dagml:DataRequirement"],
444 "dagml:requirement_key": key,
445 "dagml:node_id": requirement.node_id,
446 "dagml:input_name": requirement.input_name,
447 "dagml:schema_fingerprint": requirement.schema_fingerprint,
448 "dagml:plan_fingerprint": requirement.plan_fingerprint,
449 "dagml:relation_fingerprint": requirement.relation_fingerprint,
450 "dagml:feature_set_id": requirement.feature_set_id,
451 }),
452 );
453 }
454 for (key, envelope) in data_envelopes {
455 entity.insert(
456 data_envelope_entity_id(key),
457 json!({
458 "prov:type": ["prov:Entity", "dagml:ExternalDataPlanEnvelope"],
459 "dagml:envelope_key": key,
460 "dagml:schema_version": envelope.schema_version,
461 "dagml:schema_fingerprint": envelope.schema_fingerprint,
462 "dagml:plan_fingerprint": envelope.plan_fingerprint,
463 "dagml:relation_fingerprint": envelope.relation_fingerprint,
464 }),
465 );
466 }
467 for requirement in &bundle.prediction_requirements {
468 let key = requirement.key();
469 entity.insert(
470 prediction_requirement_entity_id(&key),
471 json!({
472 "prov:type": ["prov:Entity", "dagml:PredictionRequirement"],
473 "dagml:requirement_key": key,
474 "dagml:producer_node": requirement.producer_node,
475 "dagml:consumer_node": requirement.consumer_node,
476 "dagml:prediction_level": requirement.prediction_level,
477 "dagml:fold_ids": requirement.fold_ids,
478 "dagml:unit_ids": requirement.unit_ids,
479 "dagml:sample_ids": requirement.sample_ids,
480 "dagml:prediction_width": requirement.prediction_width,
481 "dagml:target_names": requirement.target_names,
482 }),
483 );
484 }
485 for cache in &bundle.prediction_caches {
486 entity.insert(
487 prediction_cache_entity_id(&cache.cache_id),
488 json!({
489 "prov:type": ["prov:Entity", "dagml:PredictionCache"],
490 "dagml:requirement_key": cache.requirement_key,
491 "dagml:cache_id": cache.cache_id,
492 "dagml:format": cache.format,
493 "dagml:prediction_level": cache.prediction_level,
494 "dagml:unit_ids": cache.unit_ids,
495 "dagml:block_count": cache.block_count,
496 "dagml:row_count": cache.row_count,
497 "dagml:content_fingerprint": cache.content_fingerprint,
498 }),
499 );
500 }
501 if let Some(manifest) = prediction_cache_manifest {
502 entity.insert(
503 "dagml:file:prediction-cache-manifest".to_string(),
504 json!({
505 "prov:type": ["prov:Entity", "dagml:PredictionCacheManifest"],
506 "dagml:file": FILE_PREDICTION_CACHE_MANIFEST_FILE,
507 "dagml:schema_version": manifest.schema_version,
508 "dagml:cache_count": manifest.caches.len(),
509 }),
510 );
511 }
512 for record in &bundle.refit_artifacts {
513 entity.insert(
514 artifact_entity_id(&record.artifact.id),
515 json!({
516 "prov:type": ["prov:Entity", "dagml:ModelArtifact"],
517 "dagml:artifact_id": record.artifact.id,
518 "dagml:kind": record.artifact.kind,
519 "dagml:node_id": record.node_id,
520 "dagml:controller_id": record.controller_id,
521 "dagml:backend": record.artifact.backend,
522 "dagml:uri": record.artifact.uri,
523 "dagml:content_fingerprint": record.artifact.content_fingerprint,
524 "dagml:size_bytes": record.artifact.size_bytes,
525 "dagml:plugin": record.artifact.plugin,
526 "dagml:plugin_version": record.artifact.plugin_version,
527 "dagml:params_fingerprint": record.params_fingerprint,
528 "dagml:training_loss_fingerprint": record.training_loss_fingerprint,
529 "dagml:data_requirement_keys": record.data_requirement_keys,
530 "dagml:prediction_requirement_keys": record.prediction_requirement_keys,
531 }),
532 );
533 }
534 if let Some(manifest) = artifact_manifest {
535 entity.insert(
536 "dagml:file:artifact-manifest".to_string(),
537 json!({
538 "prov:type": ["prov:Entity", "dagml:ArtifactManifest"],
539 "dagml:file": FILE_ARTIFACT_MANIFEST_FILE,
540 "dagml:schema_version": manifest.schema_version,
541 "dagml:artifact_count": manifest.artifacts.len(),
542 }),
543 );
544 }
545 for record in lineage {
546 entity.insert(
547 lineage_record_entity_id(&record.record_id),
548 json!({
549 "prov:type": ["prov:Entity", "dagml:LineageRecord"],
550 "dagml:lineage_id": record.record_id,
551 "dagml:run_id": record.run_id,
552 "dagml:node_id": record.node_id,
553 "dagml:phase": record.phase,
554 "dagml:controller_id": record.controller_id,
555 "dagml:variant_id": record.variant_id,
556 "dagml:fold_id": record.fold_id,
557 "dagml:branch_path": record.branch_path,
558 "dagml:input_lineage": record.input_lineage,
559 "dagml:artifact_refs": record
560 .artifact_refs
561 .iter()
562 .map(|artifact| artifact.id.clone())
563 .collect::<Vec<_>>(),
564 }),
565 );
566 }
567
568 let mut agent = BTreeMap::<String, Value>::new();
569 agent.insert(
570 coordinator_agent_id.clone(),
571 json!({
572 "prov:type": ["prov:Agent", "dagml:Coordinator"],
573 "dagml:name": "dag-ml",
574 "dagml:provenance_schema_version": RESEARCH_PROVENANCE_SCHEMA_VERSION,
575 }),
576 );
577 for manifest in plan.controller_manifests.values() {
578 agent.insert(
579 controller_agent_id(manifest.controller_id.as_str()),
580 json!({
581 "prov:type": ["prov:Agent", "dagml:Controller"],
582 "dagml:controller_id": manifest.controller_id,
583 "dagml:controller_version": manifest.controller_version,
584 "dagml:operator_kind": manifest.operator_kind,
585 "dagml:fit_scope": manifest.fit_scope,
586 "dagml:rng_policy": manifest.rng_policy,
587 "dagml:artifact_policy": manifest.artifact_policy,
588 "dagml:capabilities": manifest.capabilities,
589 }),
590 );
591 }
592
593 let mut activity = BTreeMap::<String, Value>::new();
594 activity.insert(
595 packaging_activity_id.clone(),
596 json!({
597 "prov:type": ["prov:Activity", "dagml:BundlePackaging"],
598 "dagml:bundle_id": bundle.bundle_id,
599 "dagml:plan_id": bundle.plan_id,
600 "dagml:selected_variant_id": bundle.selected_variant_id,
601 }),
602 );
603 for record in lineage {
604 activity.insert(
605 lineage_activity_id(record),
606 json!({
607 "prov:type": ["prov:Activity", "dagml:NodeExecution"],
608 "dagml:lineage_id": record.record_id,
609 "dagml:run_id": record.run_id,
610 "dagml:node_id": record.node_id,
611 "dagml:phase": record.phase,
612 "dagml:controller_id": record.controller_id,
613 "dagml:controller_version": record.controller_version,
614 "dagml:variant_id": record.variant_id,
615 "dagml:fold_id": record.fold_id,
616 "dagml:branch_path": record.branch_path,
617 "dagml:params_fingerprint": record.params_fingerprint,
618 "dagml:data_model_shape_fingerprint": record.data_model_shape_fingerprint,
619 "dagml:aggregation_policy_fingerprint": record.aggregation_policy_fingerprint,
620 "dagml:seed": record.seed,
621 "dagml:unsafe_flags": record.unsafe_flags,
622 "dagml:metrics": record.metrics,
623 "dagml:loss_attestations": record.loss_attestations,
624 }),
625 );
626 }
627
628 let mut used = BTreeMap::<String, Value>::new();
629 used.insert(
630 "dagml:used:bundle-plan".to_string(),
631 json!({
632 "prov:activity": packaging_activity_id,
633 "prov:entity": plan_entity_id,
634 }),
635 );
636 for record in lineage {
637 for input_id in &record.input_lineage {
638 used.insert(
639 format!("dagml:used:{}:{}", record.record_id, input_id),
640 json!({
641 "prov:activity": lineage_activity_id(record),
642 "prov:entity": lineage_record_entity_id(input_id),
643 "dagml:input_lineage_id": input_id,
644 }),
645 );
646 }
647 }
648
649 let lineage_by_artifact = lineage_artifact_index(lineage);
650 let mut was_generated_by = BTreeMap::<String, Value>::new();
651 was_generated_by.insert(
652 "dagml:generated:bundle".to_string(),
653 json!({
654 "prov:entity": bundle_entity_id,
655 "prov:activity": packaging_activity_id,
656 }),
657 );
658 for record in lineage {
659 was_generated_by.insert(
660 format!("dagml:generated:lineage:{}", record.record_id),
661 json!({
662 "prov:entity": lineage_record_entity_id(&record.record_id),
663 "prov:activity": lineage_activity_id(record),
664 }),
665 );
666 }
667 for record in &bundle.refit_artifacts {
668 let activity_id = lineage_by_artifact
669 .get(&record.artifact.id)
670 .cloned()
671 .unwrap_or_else(|| packaging_activity_id.clone());
672 was_generated_by.insert(
673 format!("dagml:generated:artifact:{}", record.artifact.id),
674 json!({
675 "prov:entity": artifact_entity_id(&record.artifact.id),
676 "prov:activity": activity_id,
677 }),
678 );
679 }
680
681 let mut was_derived_from = BTreeMap::<String, Value>::new();
682 was_derived_from.insert(
683 "dagml:derived:bundle-plan".to_string(),
684 json!({
685 "prov:generatedEntity": bundle_entity_id,
686 "prov:usedEntity": plan_entity_id,
687 }),
688 );
689 for record in &bundle.refit_artifacts {
690 for key in &record.data_requirement_keys {
691 was_derived_from.insert(
692 format!("dagml:derived:{}:data:{key}", record.artifact.id),
693 json!({
694 "prov:generatedEntity": artifact_entity_id(&record.artifact.id),
695 "prov:usedEntity": data_requirement_entity_id(key),
696 "dagml:refit_dependency": "data_requirement",
697 }),
698 );
699 }
700 for key in &record.prediction_requirement_keys {
701 was_derived_from.insert(
702 format!("dagml:derived:{}:prediction:{key}", record.artifact.id),
703 json!({
704 "prov:generatedEntity": artifact_entity_id(&record.artifact.id),
705 "prov:usedEntity": prediction_requirement_entity_id(key),
706 "dagml:refit_dependency": "prediction_requirement",
707 "dagml:oof_dependency": true,
708 }),
709 );
710 }
711 }
712 for cache in &bundle.prediction_caches {
713 was_derived_from.insert(
714 format!("dagml:derived:cache:{}", cache.cache_id),
715 json!({
716 "prov:generatedEntity": prediction_cache_entity_id(&cache.cache_id),
717 "prov:usedEntity": prediction_requirement_entity_id(&cache.requirement_key),
718 }),
719 );
720 }
721 for record in lineage {
722 for input_id in &record.input_lineage {
723 was_derived_from.insert(
724 format!("dagml:derived:lineage:{}:{input_id}", record.record_id),
725 json!({
726 "prov:generatedEntity": lineage_record_entity_id(&record.record_id),
727 "prov:usedEntity": lineage_record_entity_id(input_id),
728 "dagml:lineage_dependency": true,
729 }),
730 );
731 }
732 }
733
734 let mut was_associated_with = BTreeMap::<String, Value>::new();
735 was_associated_with.insert(
736 "dagml:associated:bundle-packaging".to_string(),
737 json!({
738 "prov:activity": packaging_activity_id,
739 "prov:agent": coordinator_agent_id,
740 }),
741 );
742 for record in lineage {
743 was_associated_with.insert(
744 format!("dagml:associated:{}", record.record_id),
745 json!({
746 "prov:activity": lineage_activity_id(record),
747 "prov:agent": controller_agent_id(record.controller_id.as_str()),
748 }),
749 );
750 }
751
752 Ok(json!({
753 "@context": {
754 "prov": "http://www.w3.org/ns/prov#",
755 "dagml": "https://dag-ml.dev/ns#",
756 },
757 "entity": entity,
758 "activity": activity,
759 "agent": agent,
760 "used": used,
761 "wasGeneratedBy": was_generated_by,
762 "wasDerivedFrom": was_derived_from,
763 "wasAssociatedWith": was_associated_with,
764 }))
765}
766
767fn build_ro_crate_metadata(
768 plan: &ExecutionPlan,
769 bundle: &ExecutionBundle,
770 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
771 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
772 artifact_manifest: Option<&FileArtifactManifest>,
773) -> Result<Value> {
774 let mut has_part = vec![
775 json!({"@id": "execution_plan.json"}),
776 json!({"@id": "execution_bundle.json"}),
777 json!({"@id": PROV_JSONLD_FILE}),
778 ];
779 let mut graph = vec![
780 json!({
781 "@id": RO_CRATE_METADATA_FILE,
782 "@type": "CreativeWork",
783 "about": {"@id": "./"},
784 "conformsTo": {"@id": "https://w3id.org/ro/crate/1.1"},
785 }),
786 json!({
787 "@id": "./",
788 "@type": "Dataset",
789 "name": format!("DAG-ML research bundle {}", bundle.bundle_id),
790 "mainEntity": {"@id": "#workflow"},
791 "hasPart": has_part.clone(),
792 "dagml:schema_version": RESEARCH_PROVENANCE_SCHEMA_VERSION,
793 "dagml:bundle_id": bundle.bundle_id,
794 "dagml:plan_id": plan.id,
795 "dagml:unsafe_flags": bundle.unsafe_flags,
796 }),
797 json!({
798 "@id": "#workflow",
799 "@type": ["ComputationalWorkflow", "SoftwareSourceCode"],
800 "name": "DAG-ML compiled workflow",
801 "programmingLanguage": "Rust",
802 "dagml:plan_id": plan.id,
803 "dagml:graph_fingerprint": plan.graph_fingerprint,
804 "dagml:campaign_fingerprint": plan.campaign_fingerprint,
805 "dagml:controller_fingerprint": plan.controller_fingerprint,
806 "dagml:selected_variant_id": bundle.selected_variant_id,
807 "dagml:variant_count": plan.variants.len(),
808 }),
809 file_entity(
810 "execution_plan.json",
811 "DAG-ML execution plan",
812 "dagml:ExecutionPlan",
813 ),
814 file_entity(
815 "execution_bundle.json",
816 "DAG-ML execution bundle",
817 "dagml:ExecutionBundle",
818 ),
819 file_entity(PROV_JSONLD_FILE, "DAG-ML W3C PROV export", "prov:Bundle"),
820 ];
821
822 if prediction_cache_manifest.is_some() {
823 has_part.push(json!({"@id": FILE_PREDICTION_CACHE_MANIFEST_FILE}));
824 graph.push(file_entity(
825 FILE_PREDICTION_CACHE_MANIFEST_FILE,
826 "DAG-ML prediction cache manifest",
827 "dagml:PredictionCacheManifest",
828 ));
829 }
830 if artifact_manifest.is_some() {
831 has_part.push(json!({"@id": FILE_ARTIFACT_MANIFEST_FILE}));
832 graph.push(file_entity(
833 FILE_ARTIFACT_MANIFEST_FILE,
834 "DAG-ML artifact manifest",
835 "dagml:ArtifactManifest",
836 ));
837 }
838 for (key, envelope) in data_envelopes {
839 let id = format!("data_envelopes/{key}.json");
840 has_part.push(json!({"@id": id}));
841 graph.push(json!({
842 "@id": id,
843 "@type": ["File", "dagml:ExternalDataPlanEnvelope"],
844 "name": format!("DAG-ML data envelope {key}"),
845 "dagml:envelope_key": key,
846 "dagml:schema_version": envelope.schema_version,
847 "dagml:schema_fingerprint": envelope.schema_fingerprint,
848 "dagml:plan_fingerprint": envelope.plan_fingerprint,
849 "dagml:relation_fingerprint": envelope.relation_fingerprint,
850 }));
851 }
852
853 graph[1]["hasPart"] = Value::Array(has_part);
854
855 for manifest in plan.controller_manifests.values() {
856 graph.push(json!({
857 "@id": controller_agent_id(manifest.controller_id.as_str()),
858 "@type": ["SoftwareApplication", "dagml:Controller"],
859 "name": manifest.controller_id,
860 "softwareVersion": manifest.controller_version,
861 "dagml:operator_kind": manifest.operator_kind,
862 "dagml:capabilities": manifest.capabilities,
863 "dagml:artifact_policy": manifest.artifact_policy,
864 }));
865 }
866 for artifact in &bundle.refit_artifacts {
867 graph.push(json!({
868 "@id": artifact_entity_id(&artifact.artifact.id),
869 "@type": ["File", "dagml:ModelArtifact"],
870 "name": artifact.artifact.id,
871 "encodingFormat": artifact.artifact.kind,
872 "dagml:node_id": artifact.node_id,
873 "dagml:controller_id": artifact.controller_id,
874 "dagml:backend": artifact.artifact.backend,
875 "dagml:uri": artifact.artifact.uri,
876 "dagml:content_fingerprint": artifact.artifact.content_fingerprint,
877 "dagml:plugin": artifact.artifact.plugin,
878 "dagml:plugin_version": artifact.artifact.plugin_version,
879 "dagml:refit_data_requirement_keys": artifact.data_requirement_keys,
880 "dagml:refit_prediction_requirement_keys": artifact.prediction_requirement_keys,
881 }));
882 }
883
884 Ok(json!({
885 "@context": [
886 "https://w3id.org/ro/crate/1.1/context",
887 {
888 "dagml": "https://dag-ml.dev/ns#",
889 "prov": "http://www.w3.org/ns/prov#",
890 }
891 ],
892 "@graph": graph,
893 }))
894}
895
896fn file_entity(id: &str, name: &str, dagml_type: &str) -> Value {
897 json!({
898 "@id": id,
899 "@type": ["File", dagml_type],
900 "name": name,
901 })
902}
903
904fn add_json_package_file<T: Serialize + ?Sized>(
905 files: &mut BTreeMap<String, ResearchProvenancePackageFile>,
906 path: &str,
907 value: &T,
908 label: &str,
909) -> Result<()> {
910 validate_package_path(path)?;
911 let mut bytes = serde_json::to_vec_pretty(value).map_err(|err| {
912 DagMlError::RuntimeValidation(format!("failed to serialize {label}: {err}"))
913 })?;
914 bytes.push(b'\n');
915 let sha256 = sha256_hex(&bytes);
916 let previous = files.insert(
917 path.to_string(),
918 ResearchProvenancePackageFile {
919 path: path.to_string(),
920 sha256,
921 size_bytes: bytes.len(),
922 bytes,
923 },
924 );
925 if previous.is_some() {
926 return Err(DagMlError::RuntimeValidation(format!(
927 "duplicate research provenance package file `{path}`"
928 )));
929 }
930 Ok(())
931}
932
933fn validate_package_path(path: &str) -> Result<()> {
934 if path.is_empty() {
935 return Err(DagMlError::RuntimeValidation(
936 "research provenance package path is empty".to_string(),
937 ));
938 }
939 if path.starts_with('/') || path.starts_with('\\') {
940 return Err(DagMlError::RuntimeValidation(format!(
941 "research provenance package path `{path}` must be relative"
942 )));
943 }
944 if path.chars().any(char::is_control) {
945 return Err(DagMlError::RuntimeValidation(format!(
946 "research provenance package path `{path}` has control characters"
947 )));
948 }
949 for segment in path.split(['/', '\\']) {
950 if segment.is_empty() || segment == "." || segment == ".." {
951 return Err(DagMlError::RuntimeValidation(format!(
952 "research provenance package path `{path}` has an invalid path component"
953 )));
954 }
955 }
956 Ok(())
957}
958
959fn data_envelope_file_path(key: &str) -> Result<String> {
960 if key.contains(['/', '\\']) {
961 return Err(DagMlError::RuntimeValidation(format!(
962 "data envelope key `{key}` cannot be used as a research provenance package path"
963 )));
964 }
965 Ok(format!("data_envelopes/{key}.json"))
966}
967
968fn require_package_file<'a>(files: &'a BTreeMap<String, Vec<u8>>, path: &str) -> Result<&'a [u8]> {
969 files.get(path).map(Vec::as_slice).ok_or_else(|| {
970 DagMlError::RuntimeValidation(format!("research provenance package is missing `{path}`"))
971 })
972}
973
974fn parse_package_json<T: DeserializeOwned>(bytes: &[u8], path: &str) -> Result<T> {
975 serde_json::from_slice(bytes).map_err(|err| {
976 DagMlError::RuntimeValidation(format!(
977 "failed to parse research provenance package JSON `{path}`: {err}"
978 ))
979 })
980}
981
982fn parse_package_data_envelopes(
983 files: &BTreeMap<String, Vec<u8>>,
984) -> Result<BTreeMap<String, ExternalDataPlanEnvelope>> {
985 let mut envelopes = BTreeMap::new();
986 for (path, bytes) in files {
987 let Some(key) = path
988 .strip_prefix("data_envelopes/")
989 .and_then(|suffix| suffix.strip_suffix(".json"))
990 else {
991 continue;
992 };
993 if key.is_empty() || key.contains(['/', '\\']) {
994 return Err(DagMlError::RuntimeValidation(format!(
995 "research provenance data envelope path `{path}` has an invalid key"
996 )));
997 }
998 let previous = envelopes.insert(key.to_string(), parse_package_json(bytes, path)?);
999 if previous.is_some() {
1000 return Err(DagMlError::RuntimeValidation(format!(
1001 "duplicate research provenance data envelope key `{key}`"
1002 )));
1003 }
1004 }
1005 Ok(envelopes)
1006}
1007
1008fn validate_prov_jsonld_root(prov_jsonld: Value) -> Result<()> {
1009 if prov_jsonld.get("@context").is_none()
1010 || prov_jsonld.get("entity").is_none()
1011 || prov_jsonld.get("activity").is_none()
1012 || prov_jsonld.get("agent").is_none()
1013 {
1014 return Err(DagMlError::RuntimeValidation(
1015 "research provenance PROV JSON-LD root is missing required sections".to_string(),
1016 ));
1017 }
1018 Ok(())
1019}
1020
1021fn validate_openlineage_namespace(namespace: &str) -> Result<()> {
1022 if namespace.trim().is_empty() {
1023 return Err(DagMlError::RuntimeValidation(
1024 "OpenLineage namespace must not be empty".to_string(),
1025 ));
1026 }
1027 if namespace.chars().any(char::is_control) {
1028 return Err(DagMlError::RuntimeValidation(
1029 "OpenLineage namespace contains control characters".to_string(),
1030 ));
1031 }
1032 Ok(())
1033}
1034
1035fn validate_openlineage_event_time(event_time: &str) -> Result<()> {
1036 if event_time.trim().is_empty() || !event_time.contains('T') {
1037 return Err(DagMlError::RuntimeValidation(
1038 "OpenLineage event_time must be a non-empty RFC3339-like timestamp".to_string(),
1039 ));
1040 }
1041 if event_time.chars().any(char::is_control) {
1042 return Err(DagMlError::RuntimeValidation(
1043 "OpenLineage event_time contains control characters".to_string(),
1044 ));
1045 }
1046 Ok(())
1047}
1048
1049fn dagml_openlineage_reproducibility_run_facet(
1050 plan: &ExecutionPlan,
1051 bundle: &ExecutionBundle,
1052) -> Value {
1053 json!({
1054 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlReproducibilityRunFacet"),
1055 "plan_id": plan.id,
1056 "bundle_id": bundle.bundle_id,
1057 "graph_fingerprint": bundle.graph_fingerprint,
1058 "campaign_fingerprint": bundle.campaign_fingerprint,
1059 "controller_fingerprint": bundle.controller_fingerprint,
1060 "selected_variant_id": bundle.selected_variant_id,
1061 "variant_count": plan.variants.len(),
1062 "unsafe_flags": bundle.unsafe_flags,
1063 })
1064}
1065
1066fn dagml_openlineage_oof_safety_run_facet(
1067 bundle: &ExecutionBundle,
1068 lineage: &[LineageRecord],
1069) -> Value {
1070 json!({
1071 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlOofSafetyRunFacet"),
1072 "prediction_requirement_count": bundle.prediction_requirements.len(),
1073 "prediction_cache_count": bundle.prediction_caches.len(),
1074 "lineage_record_count": lineage.len(),
1075 "requires_oof_prediction_count": bundle.prediction_requirements.len(),
1076 "refit_artifact_count": bundle.refit_artifacts.len(),
1077 })
1078}
1079
1080fn dagml_openlineage_plan_job_facet(plan: &ExecutionPlan, bundle: &ExecutionBundle) -> Value {
1081 json!({
1082 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlPlanJobFacet"),
1083 "plan_id": plan.id,
1084 "bundle_id": bundle.bundle_id,
1085 "node_count": plan.node_plans.len(),
1086 "controller_count": plan.controller_manifests.len(),
1087 "has_fold_set": plan.fold_set.is_some(),
1088 "selected_variant_id": bundle.selected_variant_id,
1089 })
1090}
1091
1092fn openlineage_input_datasets(
1093 bundle: &ExecutionBundle,
1094 data_envelopes: &BTreeMap<String, ExternalDataPlanEnvelope>,
1095) -> Vec<Value> {
1096 bundle
1097 .data_requirements
1098 .iter()
1099 .map(|requirement| {
1100 let key = requirement.key();
1101 let envelope = data_envelopes.get(&key);
1102 json!({
1103 "namespace": "dagml:data-requirement",
1104 "name": key,
1105 "facets": {
1106 "dagml_contract": {
1107 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1108 "node_id": requirement.node_id,
1109 "input_name": requirement.input_name,
1110 "schema_fingerprint": requirement.schema_fingerprint,
1111 "plan_fingerprint": requirement.plan_fingerprint,
1112 "relation_fingerprint": requirement.relation_fingerprint,
1113 "feature_set_id": requirement.feature_set_id,
1114 "envelope_schema_fingerprint": envelope.map(|envelope| envelope.schema_fingerprint.clone()),
1115 "envelope_plan_fingerprint": envelope.map(|envelope| envelope.plan_fingerprint.clone()),
1116 }
1117 }
1118 })
1119 })
1120 .collect()
1121}
1122
1123fn openlineage_output_datasets(
1124 bundle: &ExecutionBundle,
1125 prediction_cache_manifest: Option<&FilePredictionCacheManifest>,
1126 artifact_manifest: Option<&FileArtifactManifest>,
1127) -> Vec<Value> {
1128 let mut outputs = vec![json!({
1129 "namespace": "dagml:bundle",
1130 "name": bundle.bundle_id,
1131 "facets": {
1132 "dagml_contract": {
1133 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1134 "schema_version": bundle.schema_version,
1135 "plan_id": bundle.plan_id,
1136 "selected_variant_id": bundle.selected_variant_id,
1137 "graph_fingerprint": bundle.graph_fingerprint,
1138 "campaign_fingerprint": bundle.campaign_fingerprint,
1139 "controller_fingerprint": bundle.controller_fingerprint,
1140 }
1141 }
1142 })];
1143 for cache in &bundle.prediction_caches {
1144 outputs.push(json!({
1145 "namespace": "dagml:prediction-cache",
1146 "name": cache.cache_id,
1147 "facets": {
1148 "dagml_contract": {
1149 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1150 "requirement_key": cache.requirement_key,
1151 "prediction_level": cache.prediction_level,
1152 "row_count": cache.row_count,
1153 "block_count": cache.block_count,
1154 "content_fingerprint": cache.content_fingerprint,
1155 "has_file_manifest": prediction_cache_manifest.is_some(),
1156 }
1157 }
1158 }));
1159 }
1160 for artifact in &bundle.refit_artifacts {
1161 outputs.push(json!({
1162 "namespace": "dagml:artifact",
1163 "name": artifact.artifact.id,
1164 "facets": {
1165 "dagml_contract": {
1166 "_schemaURL": format!("{DAGML_OPENLINEAGE_FACET_SCHEMA_URL}#/$defs/DagmlDatasetContractFacet"),
1167 "node_id": artifact.node_id,
1168 "controller_id": artifact.controller_id,
1169 "backend": artifact.artifact.backend,
1170 "uri": artifact.artifact.uri,
1171 "content_fingerprint": artifact.artifact.content_fingerprint,
1172 "plugin": artifact.artifact.plugin,
1173 "plugin_version": artifact.artifact.plugin_version,
1174 "data_requirement_keys": artifact.data_requirement_keys,
1175 "prediction_requirement_keys": artifact.prediction_requirement_keys,
1176 "has_file_manifest": artifact_manifest.is_some(),
1177 }
1178 }
1179 }));
1180 }
1181 outputs
1182}
1183
1184fn openlineage_run_id(plan: &ExecutionPlan, bundle: &ExecutionBundle) -> String {
1185 let input = format!("dag-ml/openlineage/run/{}/{}", plan.id, bundle.bundle_id);
1186 let digest = Sha256::digest(input.as_bytes());
1187 let mut bytes = [0u8; 16];
1188 bytes.copy_from_slice(&digest[..16]);
1189 bytes[6] = (bytes[6] & 0x0f) | 0x50;
1190 bytes[8] = (bytes[8] & 0x3f) | 0x80;
1191 format!(
1192 "{:02x}{:02x}{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}-{:02x}{:02x}{:02x}{:02x}{:02x}{:02x}",
1193 bytes[0],
1194 bytes[1],
1195 bytes[2],
1196 bytes[3],
1197 bytes[4],
1198 bytes[5],
1199 bytes[6],
1200 bytes[7],
1201 bytes[8],
1202 bytes[9],
1203 bytes[10],
1204 bytes[11],
1205 bytes[12],
1206 bytes[13],
1207 bytes[14],
1208 bytes[15],
1209 )
1210}
1211
1212fn validate_ro_crate_package_checksums(
1213 ro_crate_metadata: &Value,
1214 files: &BTreeMap<String, Vec<u8>>,
1215) -> Result<usize> {
1216 let graph = ro_crate_metadata
1217 .get("@graph")
1218 .and_then(Value::as_array)
1219 .ok_or_else(|| {
1220 DagMlError::RuntimeValidation("RO-Crate metadata has no @graph array".to_string())
1221 })?;
1222 let root = graph
1223 .iter()
1224 .find(|entry| entry.get("@id").and_then(Value::as_str) == Some("./"))
1225 .ok_or_else(|| {
1226 DagMlError::RuntimeValidation("RO-Crate metadata has no root dataset".to_string())
1227 })?;
1228 if root.get("dagml:schema_version").and_then(Value::as_u64)
1229 != Some(RESEARCH_PROVENANCE_SCHEMA_VERSION as u64)
1230 {
1231 return Err(DagMlError::RuntimeValidation(format!(
1232 "RO-Crate root has unsupported dagml:schema_version, expected {}",
1233 RESEARCH_PROVENANCE_SCHEMA_VERSION
1234 )));
1235 }
1236 let root_has_part = root
1237 .get("hasPart")
1238 .and_then(Value::as_array)
1239 .ok_or_else(|| {
1240 DagMlError::RuntimeValidation("RO-Crate root hasPart is not an array".to_string())
1241 })?;
1242 let root_has_part_ids = root_has_part
1243 .iter()
1244 .filter_map(|entry| entry.get("@id").and_then(Value::as_str))
1245 .collect::<BTreeSet<_>>();
1246
1247 let mut checksummed = 0;
1248 for (path, bytes) in files {
1249 if path == RO_CRATE_METADATA_FILE {
1250 continue;
1251 }
1252 if !root_has_part_ids.contains(path.as_str()) {
1253 return Err(DagMlError::RuntimeValidation(format!(
1254 "RO-Crate root does not list package file `{path}` in hasPart"
1255 )));
1256 }
1257 let entry = graph
1258 .iter()
1259 .find(|entry| entry.get("@id").and_then(Value::as_str) == Some(path.as_str()))
1260 .ok_or_else(|| {
1261 DagMlError::RuntimeValidation(format!(
1262 "RO-Crate metadata is missing package file entry `{path}`"
1263 ))
1264 })?;
1265 let expected_sha256 = sha256_hex(bytes);
1266 let declared_sha256 = entry.get("sha256").and_then(Value::as_str).ok_or_else(|| {
1267 DagMlError::RuntimeValidation(format!(
1268 "RO-Crate package file `{path}` is missing sha256"
1269 ))
1270 })?;
1271 if declared_sha256 != expected_sha256 {
1272 return Err(DagMlError::RuntimeValidation(format!(
1273 "RO-Crate package file `{path}` sha256 mismatch"
1274 )));
1275 }
1276 let declared_dagml_sha256 = entry
1277 .get("dagml:sha256")
1278 .and_then(Value::as_str)
1279 .ok_or_else(|| {
1280 DagMlError::RuntimeValidation(format!(
1281 "RO-Crate package file `{path}` is missing dagml:sha256"
1282 ))
1283 })?;
1284 if declared_dagml_sha256 != declared_sha256 {
1285 return Err(DagMlError::RuntimeValidation(format!(
1286 "RO-Crate package file `{path}` has inconsistent checksum fields"
1287 )));
1288 }
1289 if entry.get("contentSize").and_then(Value::as_u64) != Some(bytes.len() as u64) {
1290 return Err(DagMlError::RuntimeValidation(format!(
1291 "RO-Crate package file `{path}` contentSize mismatch"
1292 )));
1293 }
1294 if entry.get("encodingFormat").and_then(Value::as_str) != Some("application/json") {
1295 return Err(DagMlError::RuntimeValidation(format!(
1296 "RO-Crate package file `{path}` must declare application/json"
1297 )));
1298 }
1299 checksummed += 1;
1300 }
1301
1302 for has_part_id in root_has_part_ids {
1303 if has_part_id != RO_CRATE_METADATA_FILE && !files.contains_key(has_part_id) {
1304 return Err(DagMlError::RuntimeValidation(format!(
1305 "RO-Crate root references missing package file `{has_part_id}`"
1306 )));
1307 }
1308 }
1309
1310 Ok(checksummed)
1311}
1312
1313fn annotate_ro_crate_package_files(
1314 ro_crate_metadata: &mut Value,
1315 files: &BTreeMap<String, ResearchProvenancePackageFile>,
1316) -> Result<()> {
1317 let graph = ro_crate_metadata
1318 .get_mut("@graph")
1319 .and_then(Value::as_array_mut)
1320 .ok_or_else(|| {
1321 DagMlError::RuntimeValidation("RO-Crate metadata has no @graph array".to_string())
1322 })?;
1323
1324 let mut existing_ids = graph
1325 .iter()
1326 .filter_map(|entry| entry.get("@id").and_then(Value::as_str).map(str::to_string))
1327 .collect::<BTreeSet<_>>();
1328 for file in files.values() {
1329 if !existing_ids.contains(&file.path) {
1330 graph.push(file_entity(
1331 &file.path,
1332 &format!("DAG-ML contract file {}", file.path),
1333 "dagml:ContractArtifact",
1334 ));
1335 existing_ids.insert(file.path.clone());
1336 }
1337 }
1338
1339 for entry in graph.iter_mut() {
1340 let Some(id) = entry.get("@id").and_then(Value::as_str).map(str::to_string) else {
1341 continue;
1342 };
1343 let Some(file) = files.get(id.as_str()) else {
1344 continue;
1345 };
1346 let object = entry.as_object_mut().ok_or_else(|| {
1347 DagMlError::RuntimeValidation(format!("RO-Crate graph entry `{id}` is not an object"))
1348 })?;
1349 object.insert("encodingFormat".to_string(), json!("application/json"));
1350 object.insert("contentSize".to_string(), json!(file.size_bytes));
1351 object.insert("sha256".to_string(), json!(file.sha256));
1352 object.insert("dagml:sha256".to_string(), json!(file.sha256));
1353 }
1354
1355 let root = graph
1356 .iter_mut()
1357 .find(|entry| entry.get("@id") == Some(&json!("./")))
1358 .ok_or_else(|| {
1359 DagMlError::RuntimeValidation("RO-Crate metadata has no root dataset".to_string())
1360 })?;
1361 let root_object = root.as_object_mut().ok_or_else(|| {
1362 DagMlError::RuntimeValidation("RO-Crate root dataset is not an object".to_string())
1363 })?;
1364 let has_part = root_object
1365 .entry("hasPart".to_string())
1366 .or_insert_with(|| Value::Array(Vec::new()));
1367 let has_part = has_part.as_array_mut().ok_or_else(|| {
1368 DagMlError::RuntimeValidation("RO-Crate root hasPart is not an array".to_string())
1369 })?;
1370 let mut has_part_ids = has_part
1371 .iter()
1372 .filter_map(|entry| entry.get("@id").and_then(Value::as_str).map(str::to_string))
1373 .collect::<BTreeSet<_>>();
1374 for path in files.keys() {
1375 if has_part_ids.insert(path.clone()) {
1376 has_part.push(json!({"@id": path}));
1377 }
1378 }
1379 Ok(())
1380}
1381
1382fn sha256_hex(bytes: &[u8]) -> String {
1383 let digest = Sha256::digest(bytes);
1384 let mut out = String::with_capacity(digest.len() * 2);
1385 for byte in digest {
1386 use std::fmt::Write;
1387 write!(&mut out, "{byte:02x}").expect("writing to string cannot fail");
1388 }
1389 out
1390}
1391
1392fn lineage_artifact_index(lineage: &[LineageRecord]) -> BTreeMap<ArtifactId, String> {
1393 let mut index = BTreeMap::new();
1394 for record in lineage {
1395 for artifact in &record.artifact_refs {
1396 index.insert(artifact.id.clone(), lineage_activity_id(record));
1397 }
1398 }
1399 index
1400}
1401
1402fn lineage_activity_id(record: &LineageRecord) -> String {
1403 format!("dagml:activity:{}", record.record_id)
1404}
1405
1406fn controller_agent_id(controller_id: &str) -> String {
1407 format!("dagml:controller:{controller_id}")
1408}
1409
1410fn artifact_entity_id(artifact_id: &ArtifactId) -> String {
1411 format!("dagml:artifact:{artifact_id}")
1412}
1413
1414fn lineage_record_entity_id(lineage_id: &LineageId) -> String {
1415 format!("dagml:lineage-record:{lineage_id}")
1416}
1417
1418fn data_requirement_entity_id(key: &str) -> String {
1419 format!("dagml:data-requirement:{key}")
1420}
1421
1422fn data_envelope_entity_id(key: &str) -> String {
1423 format!("dagml:data-envelope:{key}")
1424}
1425
1426fn prediction_requirement_entity_id(key: &str) -> String {
1427 format!("dagml:prediction-requirement:{key}")
1428}
1429
1430fn prediction_cache_entity_id(cache_id: &str) -> String {
1431 format!("dagml:prediction-cache:{cache_id}")
1432}
1433
1434#[cfg(test)]
1435mod tests {
1436 use super::*;
1437 use crate::controller::{ControllerManifest, ControllerRegistry};
1438 use crate::ids::{ControllerId, LineageId, NodeId, RunId};
1439 use crate::plan::build_execution_plan;
1440 use crate::{CampaignSpec, GraphSpec, Phase};
1441
1442 fn fixture_plan() -> ExecutionPlan {
1443 let graph: GraphSpec =
1444 serde_json::from_str(include_str!("../../../examples/minimal_graph.json")).unwrap();
1445 let campaign: CampaignSpec = serde_json::from_str(include_str!(
1446 "../../../examples/campaign_oof_generation.json"
1447 ))
1448 .unwrap();
1449 let manifests: Vec<ControllerManifest> =
1450 serde_json::from_str(include_str!("../../../examples/controller_manifests.json"))
1451 .unwrap();
1452 let mut registry = ControllerRegistry::new();
1453 for manifest in manifests {
1454 registry.register(manifest).unwrap();
1455 }
1456 build_execution_plan("plan:cli.bundle", graph, campaign, ®istry).unwrap()
1457 }
1458
1459 fn fixture_bundle() -> ExecutionBundle {
1460 serde_json::from_str(include_str!(
1461 "../../../examples/generated/execution_bundle_minimal.json"
1462 ))
1463 .unwrap()
1464 }
1465
1466 fn fixture_lineage(plan: &ExecutionPlan) -> LineageRecord {
1467 let node_id = NodeId::new("model:base").unwrap();
1468 let node_plan = plan.node_plans.get(&node_id).unwrap();
1469 LineageRecord {
1470 record_id: LineageId::new("lineage:test:model:base").unwrap(),
1471 run_id: RunId::new("run:provenance").unwrap(),
1472 node_id,
1473 phase: Phase::Refit,
1474 controller_id: node_plan.controller_id.clone(),
1475 controller_version: node_plan.controller_version.clone(),
1476 variant_id: plan
1477 .variants
1478 .first()
1479 .map(|variant| variant.variant_id.clone()),
1480 fold_id: None,
1481 branch_path: Vec::new(),
1482 input_lineage: Vec::new(),
1483 artifact_refs: vec![fixture_bundle().refit_artifacts[0].artifact.clone()],
1484 params_fingerprint: node_plan.params_fingerprint.clone(),
1485 data_model_shape_fingerprint: None,
1486 aggregation_policy_fingerprint: None,
1487 seed: Some(42),
1488 unsafe_flags: BTreeSet::new(),
1489 metrics: BTreeMap::new(),
1490 loss_attestations: Vec::new(),
1491 early_stopping_records: Vec::new(),
1492 }
1493 }
1494
1495 #[test]
1496 fn research_provenance_export_contains_prov_and_ro_crate_contracts() {
1497 let plan = fixture_plan();
1498 let bundle = fixture_bundle();
1499 let lineage = vec![fixture_lineage(&plan)];
1500 let export = build_research_provenance_export(
1501 &plan,
1502 &bundle,
1503 &lineage,
1504 &BTreeMap::new(),
1505 None,
1506 None,
1507 )
1508 .unwrap();
1509
1510 assert_eq!(export.schema_version, RESEARCH_PROVENANCE_SCHEMA_VERSION);
1511 assert!(export.prov_jsonld["@context"]["prov"]
1512 .as_str()
1513 .unwrap()
1514 .contains("prov"));
1515 assert!(export.prov_jsonld["activity"]
1516 .as_object()
1517 .unwrap()
1518 .contains_key("dagml:activity:lineage:test:model:base"));
1519 assert!(export.prov_jsonld["agent"]
1520 .as_object()
1521 .unwrap()
1522 .contains_key("dagml:controller:controller:model.mock"));
1523 assert!(export.prov_jsonld["entity"]
1524 .as_object()
1525 .unwrap()
1526 .contains_key("dagml:artifact:artifact:model:base:refit"));
1527
1528 let graph = export.ro_crate_metadata["@graph"].as_array().unwrap();
1529 assert!(graph
1530 .iter()
1531 .any(|entry| entry["@type"].to_string().contains("ComputationalWorkflow")));
1532 assert!(graph
1533 .iter()
1534 .any(|entry| entry["@id"] == json!("lineage.prov.jsonld")));
1535 assert!(graph
1536 .iter()
1537 .any(|entry| entry["@id"] == json!("execution_bundle.json")));
1538 }
1539
1540 #[test]
1541 fn research_provenance_package_contains_contract_files_and_checksums() {
1542 let plan = fixture_plan();
1543 let bundle = fixture_bundle();
1544 let lineage = vec![fixture_lineage(&plan)];
1545 let package = build_research_provenance_package(
1546 &plan,
1547 &bundle,
1548 &lineage,
1549 &BTreeMap::new(),
1550 None,
1551 None,
1552 )
1553 .unwrap();
1554
1555 for path in [
1556 EXECUTION_PLAN_FILE,
1557 EXECUTION_BUNDLE_FILE,
1558 LINEAGE_RECORDS_FILE,
1559 PROV_JSONLD_FILE,
1560 RO_CRATE_METADATA_FILE,
1561 ] {
1562 assert!(
1563 package.files.contains_key(path),
1564 "package is missing {path}"
1565 );
1566 }
1567 for (path, file) in &package.files {
1568 assert_eq!(file.path, *path);
1569 assert_eq!(file.sha256.len(), 64, "invalid sha256 for {path}");
1570 assert!(file.size_bytes > 0, "empty package file {path}");
1571 assert_eq!(file.size_bytes, file.bytes.len());
1572 }
1573
1574 let ro_crate_file = package.files.get(RO_CRATE_METADATA_FILE).unwrap();
1575 let ro_crate_metadata: Value = serde_json::from_slice(&ro_crate_file.bytes).unwrap();
1576 let graph = ro_crate_metadata["@graph"].as_array().unwrap();
1577 for path in [
1578 EXECUTION_PLAN_FILE,
1579 EXECUTION_BUNDLE_FILE,
1580 LINEAGE_RECORDS_FILE,
1581 PROV_JSONLD_FILE,
1582 ] {
1583 let entry = graph
1584 .iter()
1585 .find(|entry| entry["@id"] == json!(path))
1586 .unwrap_or_else(|| panic!("RO-Crate metadata is missing file entry {path}"));
1587 assert_eq!(entry["sha256"].as_str().map(str::len), Some(64));
1588 assert_eq!(entry["dagml:sha256"].as_str(), entry["sha256"].as_str());
1589 assert_eq!(entry["encodingFormat"], json!("application/json"));
1590 assert!(entry["contentSize"].as_u64().unwrap() > 0);
1591 }
1592 let root = graph
1593 .iter()
1594 .find(|entry| entry["@id"] == json!("./"))
1595 .expect("RO-Crate root dataset is present");
1596 let has_part = root["hasPart"].as_array().unwrap();
1597 assert!(has_part
1598 .iter()
1599 .any(|entry| entry["@id"] == json!(LINEAGE_RECORDS_FILE)));
1600 }
1601
1602 #[test]
1603 fn research_provenance_package_validation_reopens_exported_contracts() {
1604 let plan = fixture_plan();
1605 let bundle = fixture_bundle();
1606 let lineage = vec![fixture_lineage(&plan)];
1607 let package = build_research_provenance_package(
1608 &plan,
1609 &bundle,
1610 &lineage,
1611 &BTreeMap::new(),
1612 None,
1613 None,
1614 )
1615 .unwrap();
1616 let files = package
1617 .files
1618 .iter()
1619 .map(|(path, file)| (path.clone(), file.bytes.clone()))
1620 .collect::<BTreeMap<_, _>>();
1621
1622 let validation = validate_research_provenance_package_files(&files).unwrap();
1623
1624 assert_eq!(
1625 validation.schema_version,
1626 RESEARCH_PROVENANCE_SCHEMA_VERSION
1627 );
1628 assert_eq!(validation.plan_id, plan.id.to_string());
1629 assert_eq!(validation.bundle_id, bundle.bundle_id.to_string());
1630 assert_eq!(validation.file_count, package.files.len());
1631 assert_eq!(validation.checksummed_file_count, package.files.len() - 1);
1632 assert_eq!(validation.lineage_record_count, 1);
1633 }
1634
1635 #[test]
1636 fn research_provenance_package_validation_refuses_tampered_file() {
1637 let plan = fixture_plan();
1638 let bundle = fixture_bundle();
1639 let package =
1640 build_research_provenance_package(&plan, &bundle, &[], &BTreeMap::new(), None, None)
1641 .unwrap();
1642 let mut files = package
1643 .files
1644 .iter()
1645 .map(|(path, file)| (path.clone(), file.bytes.clone()))
1646 .collect::<BTreeMap<_, _>>();
1647 files.insert(LINEAGE_RECORDS_FILE.to_string(), b"[]\n\n".to_vec());
1648
1649 let error = validate_research_provenance_package_files(&files)
1650 .unwrap_err()
1651 .to_string();
1652
1653 assert!(
1654 error.contains("sha256 mismatch"),
1655 "unexpected error: {error}"
1656 );
1657 }
1658
1659 #[test]
1660 fn openlineage_export_uses_validated_research_package() {
1661 let plan = fixture_plan();
1662 let bundle = fixture_bundle();
1663 let lineage = vec![fixture_lineage(&plan)];
1664 let package = build_research_provenance_package(
1665 &plan,
1666 &bundle,
1667 &lineage,
1668 &BTreeMap::new(),
1669 None,
1670 None,
1671 )
1672 .unwrap();
1673 let files = package
1674 .files
1675 .iter()
1676 .map(|(path, file)| (path.clone(), file.bytes.clone()))
1677 .collect::<BTreeMap<_, _>>();
1678
1679 let event = build_openlineage_run_event_from_package_files(
1680 &files,
1681 "dag-ml-test",
1682 "2026-05-27T00:00:00Z",
1683 )
1684 .unwrap();
1685
1686 assert_eq!(event["eventType"], json!("COMPLETE"));
1687 assert_eq!(event["schemaURL"], json!(OPENLINEAGE_RUN_EVENT_SCHEMA_URL));
1688 assert_eq!(event["job"]["namespace"], json!("dag-ml-test"));
1689 assert_eq!(event["run"]["runId"].as_str().map(str::len), Some(36));
1690 assert!(
1691 event["run"]["facets"]["dagml_reproducibility"]["graph_fingerprint"]
1692 .as_str()
1693 .is_some()
1694 );
1695 assert!(!event["inputs"].as_array().unwrap().is_empty());
1696 assert!(event["outputs"]
1697 .as_array()
1698 .unwrap()
1699 .iter()
1700 .any(|output| output["namespace"] == json!("dagml:bundle")));
1701 }
1702
1703 #[test]
1704 fn research_provenance_export_refuses_unknown_lineage_node() {
1705 let plan = fixture_plan();
1706 let bundle = fixture_bundle();
1707 let mut lineage = fixture_lineage(&plan);
1708 lineage.node_id = NodeId::new("model:missing").unwrap();
1709
1710 let error = build_research_provenance_export(
1711 &plan,
1712 &bundle,
1713 &[lineage],
1714 &BTreeMap::new(),
1715 None,
1716 None,
1717 )
1718 .unwrap_err()
1719 .to_string();
1720
1721 assert!(error.contains("unknown node"), "unexpected error: {error}");
1722 }
1723
1724 #[test]
1725 fn research_provenance_export_refuses_mismatched_artifact_manifest() {
1726 let plan = fixture_plan();
1727 let bundle = fixture_bundle();
1728 let mut manifest = FileArtifactManifest {
1729 bundle_id: bundle.bundle_id.clone(),
1730 schema_version: crate::runtime::FILE_ARTIFACT_MANIFEST_SCHEMA_VERSION,
1731 artifacts: Vec::new(),
1732 };
1733 manifest.bundle_id = crate::ids::BundleId::new("bundle:wrong").unwrap();
1734
1735 let error = build_research_provenance_export(
1736 &plan,
1737 &bundle,
1738 &[],
1739 &BTreeMap::new(),
1740 None,
1741 Some(&manifest),
1742 )
1743 .unwrap_err()
1744 .to_string();
1745
1746 assert!(
1747 error.contains("does not match bundle"),
1748 "unexpected error: {error}"
1749 );
1750 }
1751
1752 #[test]
1753 fn research_provenance_export_refuses_unknown_lineage_controller() {
1754 let plan = fixture_plan();
1755 let bundle = fixture_bundle();
1756 let mut lineage = fixture_lineage(&plan);
1757 lineage.controller_id = ControllerId::new("controller:missing").unwrap();
1758
1759 let error = build_research_provenance_export(
1760 &plan,
1761 &bundle,
1762 &[lineage],
1763 &BTreeMap::new(),
1764 None,
1765 None,
1766 )
1767 .unwrap_err()
1768 .to_string();
1769
1770 assert!(
1771 error.contains("unknown controller"),
1772 "unexpected error: {error}"
1773 );
1774 }
1775}