1use std::path::Path;
23
24use crate::engine::Engine;
25use crate::pipeline::{Facet, Ingest, Medium, Projection};
26use crate::pipeline_store::{self, PipelineConfigs};
27use crate::workspace_store::StoreError;
28
29#[derive(Debug, thiserror::Error)]
33pub enum PipelineEditError {
34 #[error("engine has no workspace root — pipeline edits require a workspace-backed engine")]
38 NoWorkspaceRoot,
39 #[error("{primitive} '{key}' already exists")]
41 AlreadyExists {
42 primitive: &'static str,
43 key: String,
44 },
45 #[error("{primitive} '{key}' does not exist")]
47 NotFound {
48 primitive: &'static str,
49 key: String,
50 },
51 #[error("{primitive} '{key}' is referenced by {referrers:?} — remove or repoint them first")]
53 Referenced {
54 primitive: &'static str,
55 key: String,
56 referrers: Vec<String>,
57 },
58 #[error("rename target {primitive} '{key}' already exists")]
60 RenameTargetExists {
61 primitive: &'static str,
62 key: String,
63 },
64 #[error("invalid {primitive} JSON: {message}")]
67 InvalidJson {
68 primitive: &'static str,
69 message: String,
70 },
71 #[error(transparent)]
73 Store(#[from] StoreError),
74}
75
76fn key(mem: &str, name: &str) -> String {
77 format!("{mem}/{name}")
78}
79
80fn medium_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
81 c.mediums.iter().any(|r| r.mem == mem && r.name == name)
82}
83
84fn facet_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
85 c.facets.iter().any(|r| r.mem == mem && r.name == name)
86}
87
88fn projection_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
89 c.projections.iter().any(|r| r.mem == mem && r.name == name)
90}
91
92fn facets_referencing_medium(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
94 c.facets
95 .iter()
96 .filter(|r| r.mem == mem && r.config.medium == name)
97 .map(|r| r.name.clone())
98 .collect()
99}
100
101fn projections_referencing_facet(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
103 c.projections
104 .iter()
105 .filter(|r| r.mem == mem && r.config.source_facets.iter().any(|f| f == name))
106 .map(|r| r.name.clone())
107 .collect()
108}
109
110fn ingests_referencing_projection(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
112 let target = key(mem, name);
113 c.ingests
114 .iter()
115 .filter(|r| r.config.projection == target)
116 .map(|r| r.name.clone())
117 .collect()
118}
119
120pub fn add_medium(
124 root: &Path,
125 mem: &str,
126 name: &str,
127 medium: &Medium,
128) -> Result<(), PipelineEditError> {
129 let configs = pipeline_store::load_pipeline_configs(root)?;
130 if medium_exists(&configs, mem, name) {
131 return Err(PipelineEditError::AlreadyExists {
132 primitive: "medium",
133 key: key(mem, name),
134 });
135 }
136 pipeline_store::write_medium(root, mem, name, medium)?;
137 Ok(())
138}
139
140pub fn update_medium(
142 root: &Path,
143 mem: &str,
144 name: &str,
145 medium: &Medium,
146) -> Result<(), PipelineEditError> {
147 let configs = pipeline_store::load_pipeline_configs(root)?;
148 if !medium_exists(&configs, mem, name) {
149 return Err(PipelineEditError::NotFound {
150 primitive: "medium",
151 key: key(mem, name),
152 });
153 }
154 pipeline_store::write_medium(root, mem, name, medium)?;
155 Ok(())
156}
157
158pub fn delete_medium(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
160 let configs = pipeline_store::load_pipeline_configs(root)?;
161 if !medium_exists(&configs, mem, name) {
162 return Err(PipelineEditError::NotFound {
163 primitive: "medium",
164 key: key(mem, name),
165 });
166 }
167 let referrers = facets_referencing_medium(&configs, mem, name);
168 if !referrers.is_empty() {
169 return Err(PipelineEditError::Referenced {
170 primitive: "medium",
171 key: key(mem, name),
172 referrers,
173 });
174 }
175 pipeline_store::delete_medium(root, mem, name)?;
176 Ok(())
177}
178
179pub fn rename_medium(
182 root: &Path,
183 mem: &str,
184 old: &str,
185 new: &str,
186) -> Result<(), PipelineEditError> {
187 if old == new {
188 return Ok(());
189 }
190 let configs = pipeline_store::load_pipeline_configs(root)?;
191 let existing = configs
192 .mediums
193 .iter()
194 .find(|r| r.mem == mem && r.name == old)
195 .ok_or_else(|| PipelineEditError::NotFound {
196 primitive: "medium",
197 key: key(mem, old),
198 })?;
199 if medium_exists(&configs, mem, new) {
200 return Err(PipelineEditError::RenameTargetExists {
201 primitive: "medium",
202 key: key(mem, new),
203 });
204 }
205 let mut renamed = existing.config.clone();
206 renamed.name = new.to_string();
207 pipeline_store::write_medium(root, mem, new, &renamed)?;
210 for facet in facets_referencing_medium(&configs, mem, old) {
211 if let Some(rec) = configs
212 .facets
213 .iter()
214 .find(|r| r.mem == mem && r.name == facet)
215 {
216 let mut updated = rec.config.clone();
217 updated.medium = new.to_string();
218 pipeline_store::write_facet(root, mem, &facet, &updated)?;
219 }
220 }
221 pipeline_store::delete_medium(root, mem, old)?;
222 Ok(())
223}
224
225pub fn add_facet(
229 root: &Path,
230 mem: &str,
231 name: &str,
232 facet: &Facet,
233) -> Result<(), PipelineEditError> {
234 let configs = pipeline_store::load_pipeline_configs(root)?;
235 if facet_exists(&configs, mem, name) {
236 return Err(PipelineEditError::AlreadyExists {
237 primitive: "facet",
238 key: key(mem, name),
239 });
240 }
241 pipeline_store::write_facet(root, mem, name, facet)?;
242 Ok(())
243}
244
245pub fn update_facet(
247 root: &Path,
248 mem: &str,
249 name: &str,
250 facet: &Facet,
251) -> Result<(), PipelineEditError> {
252 let configs = pipeline_store::load_pipeline_configs(root)?;
253 if !facet_exists(&configs, mem, name) {
254 return Err(PipelineEditError::NotFound {
255 primitive: "facet",
256 key: key(mem, name),
257 });
258 }
259 pipeline_store::write_facet(root, mem, name, facet)?;
260 Ok(())
261}
262
263pub fn delete_facet(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
265 let configs = pipeline_store::load_pipeline_configs(root)?;
266 if !facet_exists(&configs, mem, name) {
267 return Err(PipelineEditError::NotFound {
268 primitive: "facet",
269 key: key(mem, name),
270 });
271 }
272 let referrers = projections_referencing_facet(&configs, mem, name);
273 if !referrers.is_empty() {
274 return Err(PipelineEditError::Referenced {
275 primitive: "facet",
276 key: key(mem, name),
277 referrers,
278 });
279 }
280 pipeline_store::delete_facet(root, mem, name)?;
281 Ok(())
282}
283
284pub fn rename_facet(root: &Path, mem: &str, old: &str, new: &str) -> Result<(), PipelineEditError> {
287 if old == new {
288 return Ok(());
289 }
290 let configs = pipeline_store::load_pipeline_configs(root)?;
291 let existing = configs
292 .facets
293 .iter()
294 .find(|r| r.mem == mem && r.name == old)
295 .ok_or_else(|| PipelineEditError::NotFound {
296 primitive: "facet",
297 key: key(mem, old),
298 })?;
299 if facet_exists(&configs, mem, new) {
300 return Err(PipelineEditError::RenameTargetExists {
301 primitive: "facet",
302 key: key(mem, new),
303 });
304 }
305 let mut renamed = existing.config.clone();
306 renamed.name = new.to_string();
307 pipeline_store::write_facet(root, mem, new, &renamed)?;
308 for proj in projections_referencing_facet(&configs, mem, old) {
309 if let Some(rec) = configs
310 .projections
311 .iter()
312 .find(|r| r.mem == mem && r.name == proj)
313 {
314 let mut updated = rec.config.clone();
315 for f in updated.source_facets.iter_mut() {
316 if f == old {
317 *f = new.to_string();
318 }
319 }
320 pipeline_store::write_projection(root, mem, &proj, &updated)?;
321 }
322 }
323 pipeline_store::delete_facet(root, mem, old)?;
324 Ok(())
325}
326
327pub fn add_projection(
331 root: &Path,
332 mem: &str,
333 name: &str,
334 projection: &Projection,
335) -> Result<(), PipelineEditError> {
336 let configs = pipeline_store::load_pipeline_configs(root)?;
337 if projection_exists(&configs, mem, name) {
338 return Err(PipelineEditError::AlreadyExists {
339 primitive: "projection",
340 key: key(mem, name),
341 });
342 }
343 pipeline_store::write_projection(root, mem, name, projection)?;
344 Ok(())
345}
346
347pub fn update_projection(
349 root: &Path,
350 mem: &str,
351 name: &str,
352 projection: &Projection,
353) -> Result<(), PipelineEditError> {
354 let configs = pipeline_store::load_pipeline_configs(root)?;
355 if !projection_exists(&configs, mem, name) {
356 return Err(PipelineEditError::NotFound {
357 primitive: "projection",
358 key: key(mem, name),
359 });
360 }
361 pipeline_store::write_projection(root, mem, name, projection)?;
362 Ok(())
363}
364
365pub fn delete_projection(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
367 let configs = pipeline_store::load_pipeline_configs(root)?;
368 if !projection_exists(&configs, mem, name) {
369 return Err(PipelineEditError::NotFound {
370 primitive: "projection",
371 key: key(mem, name),
372 });
373 }
374 let referrers = ingests_referencing_projection(&configs, mem, name);
375 if !referrers.is_empty() {
376 return Err(PipelineEditError::Referenced {
377 primitive: "projection",
378 key: key(mem, name),
379 referrers,
380 });
381 }
382 pipeline_store::delete_projection(root, mem, name)?;
383 Ok(())
384}
385
386pub fn rename_projection(
390 root: &Path,
391 mem: &str,
392 old: &str,
393 new: &str,
394) -> Result<(), PipelineEditError> {
395 if old == new {
396 return Ok(());
397 }
398 let configs = pipeline_store::load_pipeline_configs(root)?;
399 if !projection_exists(&configs, mem, old) {
400 return Err(PipelineEditError::NotFound {
401 primitive: "projection",
402 key: key(mem, old),
403 });
404 }
405 if projection_exists(&configs, mem, new) {
406 return Err(PipelineEditError::RenameTargetExists {
407 primitive: "projection",
408 key: key(mem, new),
409 });
410 }
411 pipeline_store::rename_projection(root, mem, old, new)?;
412 let new_ref = key(mem, new);
413 for ingest in ingests_referencing_projection(&configs, mem, old) {
414 if let Some(rec) = configs.ingests.iter().find(|r| r.name == ingest) {
415 let mut updated = rec.config.clone();
416 updated.projection = new_ref.clone();
417 pipeline_store::write_ingest(root, &ingest, &updated)?;
418 }
419 }
420 Ok(())
421}
422
423fn ingest_exists(c: &PipelineConfigs, name: &str) -> bool {
432 c.ingests.iter().any(|r| r.name == name)
433}
434
435pub fn add_ingest(root: &Path, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
437 let configs = pipeline_store::load_pipeline_configs(root)?;
438 if ingest_exists(&configs, name) {
439 return Err(PipelineEditError::AlreadyExists {
440 primitive: "ingest",
441 key: name.to_string(),
442 });
443 }
444 pipeline_store::write_ingest(root, name, ingest)?;
445 Ok(())
446}
447
448pub fn update_ingest(root: &Path, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
450 let configs = pipeline_store::load_pipeline_configs(root)?;
451 if !ingest_exists(&configs, name) {
452 return Err(PipelineEditError::NotFound {
453 primitive: "ingest",
454 key: name.to_string(),
455 });
456 }
457 pipeline_store::write_ingest(root, name, ingest)?;
458 Ok(())
459}
460
461pub fn delete_ingest(root: &Path, name: &str) -> Result<(), PipelineEditError> {
463 let configs = pipeline_store::load_pipeline_configs(root)?;
464 if !ingest_exists(&configs, name) {
465 return Err(PipelineEditError::NotFound {
466 primitive: "ingest",
467 key: name.to_string(),
468 });
469 }
470 pipeline_store::delete_ingest(root, name)?;
471 Ok(())
472}
473
474pub fn rename_ingest(root: &Path, old: &str, new: &str) -> Result<(), PipelineEditError> {
477 if old == new {
478 return Ok(());
479 }
480 let configs = pipeline_store::load_pipeline_configs(root)?;
481 if !ingest_exists(&configs, old) {
482 return Err(PipelineEditError::NotFound {
483 primitive: "ingest",
484 key: old.to_string(),
485 });
486 }
487 if ingest_exists(&configs, new) {
488 return Err(PipelineEditError::RenameTargetExists {
489 primitive: "ingest",
490 key: new.to_string(),
491 });
492 }
493 pipeline_store::rename_ingest(root, old, new)?;
494 Ok(())
495}
496
497impl Engine {
506 fn pipeline_edit_root(&self) -> Result<std::path::PathBuf, PipelineEditError> {
507 self.workspace_root()
508 .map(Path::to_path_buf)
509 .ok_or(PipelineEditError::NoWorkspaceRoot)
510 }
511
512 fn refresh_pipeline_configs(&mut self, root: &Path) -> Result<(), PipelineEditError> {
513 self.set_pipeline_configs(pipeline_store::load_pipeline_configs(root)?);
514 Ok(())
515 }
516
517 pub fn add_medium(
519 &mut self,
520 mem: &str,
521 name: &str,
522 medium: &Medium,
523 ) -> Result<(), PipelineEditError> {
524 let root = self.pipeline_edit_root()?;
525 add_medium(&root, mem, name, medium)?;
526 self.refresh_pipeline_configs(&root)
527 }
528
529 pub fn update_medium(
531 &mut self,
532 mem: &str,
533 name: &str,
534 medium: &Medium,
535 ) -> Result<(), PipelineEditError> {
536 let root = self.pipeline_edit_root()?;
537 update_medium(&root, mem, name, medium)?;
538 self.refresh_pipeline_configs(&root)
539 }
540
541 pub fn delete_medium(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
543 let root = self.pipeline_edit_root()?;
544 delete_medium(&root, mem, name)?;
545 self.refresh_pipeline_configs(&root)
546 }
547
548 pub fn rename_medium(
550 &mut self,
551 mem: &str,
552 old: &str,
553 new: &str,
554 ) -> Result<(), PipelineEditError> {
555 let root = self.pipeline_edit_root()?;
556 rename_medium(&root, mem, old, new)?;
557 self.refresh_pipeline_configs(&root)
558 }
559
560 pub fn add_facet(
562 &mut self,
563 mem: &str,
564 name: &str,
565 facet: &Facet,
566 ) -> Result<(), PipelineEditError> {
567 let root = self.pipeline_edit_root()?;
568 add_facet(&root, mem, name, facet)?;
569 self.refresh_pipeline_configs(&root)
570 }
571
572 pub fn update_facet(
574 &mut self,
575 mem: &str,
576 name: &str,
577 facet: &Facet,
578 ) -> Result<(), PipelineEditError> {
579 let root = self.pipeline_edit_root()?;
580 update_facet(&root, mem, name, facet)?;
581 self.refresh_pipeline_configs(&root)
582 }
583
584 pub fn delete_facet(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
586 let root = self.pipeline_edit_root()?;
587 delete_facet(&root, mem, name)?;
588 self.refresh_pipeline_configs(&root)
589 }
590
591 pub fn rename_facet(
593 &mut self,
594 mem: &str,
595 old: &str,
596 new: &str,
597 ) -> Result<(), PipelineEditError> {
598 let root = self.pipeline_edit_root()?;
599 rename_facet(&root, mem, old, new)?;
600 self.refresh_pipeline_configs(&root)
601 }
602
603 pub fn add_projection(
605 &mut self,
606 mem: &str,
607 name: &str,
608 projection: &Projection,
609 ) -> Result<(), PipelineEditError> {
610 let root = self.pipeline_edit_root()?;
611 add_projection(&root, mem, name, projection)?;
612 self.refresh_pipeline_configs(&root)
613 }
614
615 pub fn update_projection(
617 &mut self,
618 mem: &str,
619 name: &str,
620 projection: &Projection,
621 ) -> Result<(), PipelineEditError> {
622 let root = self.pipeline_edit_root()?;
623 update_projection(&root, mem, name, projection)?;
624 self.refresh_pipeline_configs(&root)
625 }
626
627 pub fn delete_projection(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
629 let root = self.pipeline_edit_root()?;
630 delete_projection(&root, mem, name)?;
631 self.refresh_pipeline_configs(&root)
632 }
633
634 pub fn rename_projection(
636 &mut self,
637 mem: &str,
638 old: &str,
639 new: &str,
640 ) -> Result<(), PipelineEditError> {
641 let root = self.pipeline_edit_root()?;
642 rename_projection(&root, mem, old, new)?;
643 self.refresh_pipeline_configs(&root)
644 }
645
646 pub fn add_ingest(&mut self, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
648 let root = self.pipeline_edit_root()?;
649 add_ingest(&root, name, ingest)?;
650 self.refresh_pipeline_configs(&root)
651 }
652
653 pub fn update_ingest(&mut self, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
655 let root = self.pipeline_edit_root()?;
656 update_ingest(&root, name, ingest)?;
657 self.refresh_pipeline_configs(&root)
658 }
659
660 pub fn delete_ingest(&mut self, name: &str) -> Result<(), PipelineEditError> {
662 let root = self.pipeline_edit_root()?;
663 delete_ingest(&root, name)?;
664 self.refresh_pipeline_configs(&root)
665 }
666
667 pub fn rename_ingest(&mut self, old: &str, new: &str) -> Result<(), PipelineEditError> {
669 let root = self.pipeline_edit_root()?;
670 rename_ingest(&root, old, new)?;
671 self.refresh_pipeline_configs(&root)
672 }
673
674 pub fn add_medium_json(
684 &mut self,
685 mem: &str,
686 name: &str,
687 medium_json: &str,
688 ) -> Result<(), PipelineEditError> {
689 self.add_medium(mem, name, &parse_json(medium_json, "medium")?)
690 }
691
692 pub fn update_medium_json(
694 &mut self,
695 mem: &str,
696 name: &str,
697 medium_json: &str,
698 ) -> Result<(), PipelineEditError> {
699 self.update_medium(mem, name, &parse_json(medium_json, "medium")?)
700 }
701
702 pub fn add_facet_json(
704 &mut self,
705 mem: &str,
706 name: &str,
707 facet_json: &str,
708 ) -> Result<(), PipelineEditError> {
709 self.add_facet(mem, name, &parse_json(facet_json, "facet")?)
710 }
711
712 pub fn update_facet_json(
714 &mut self,
715 mem: &str,
716 name: &str,
717 facet_json: &str,
718 ) -> Result<(), PipelineEditError> {
719 self.update_facet(mem, name, &parse_json(facet_json, "facet")?)
720 }
721
722 pub fn add_projection_json(
724 &mut self,
725 mem: &str,
726 name: &str,
727 projection_json: &str,
728 ) -> Result<(), PipelineEditError> {
729 self.add_projection(mem, name, &parse_json(projection_json, "projection")?)
730 }
731
732 pub fn update_projection_json(
734 &mut self,
735 mem: &str,
736 name: &str,
737 projection_json: &str,
738 ) -> Result<(), PipelineEditError> {
739 self.update_projection(mem, name, &parse_json(projection_json, "projection")?)
740 }
741
742 pub fn add_ingest_json(
744 &mut self,
745 name: &str,
746 ingest_json: &str,
747 ) -> Result<(), PipelineEditError> {
748 self.add_ingest(name, &parse_json(ingest_json, "ingest")?)
749 }
750
751 pub fn update_ingest_json(
753 &mut self,
754 name: &str,
755 ingest_json: &str,
756 ) -> Result<(), PipelineEditError> {
757 self.update_ingest(name, &parse_json(ingest_json, "ingest")?)
758 }
759}
760
761fn parse_json<T: serde::de::DeserializeOwned>(
764 json: &str,
765 primitive: &'static str,
766) -> Result<T, PipelineEditError> {
767 serde_json::from_str(json).map_err(|e| PipelineEditError::InvalidJson {
768 primitive,
769 message: e.to_string(),
770 })
771}
772
773#[cfg(test)]
774mod tests {
775 use super::*;
776 use crate::pipeline::{
777 Ingest, IngestMode, IngestTrigger, MediumType, PatternEntry, PatternMode,
778 };
779 use tempfile::TempDir;
780
781 fn medium(name: &str) -> Medium {
782 Medium {
783 name: name.to_string(),
784 medium_type: MediumType::Codebase,
785 pointer: "../src".to_string(),
786 }
787 }
788
789 fn facet(name: &str, medium: &str) -> Facet {
790 Facet {
791 name: name.to_string(),
792 medium: medium.to_string(),
793 scope: vec![PatternEntry {
794 path: "**/*.rs".to_string(),
795 mode: PatternMode::Allow,
796 }],
797 engagement: None,
798 preparation: None,
799 }
800 }
801
802 fn projection(facets: &[&str]) -> Projection {
803 Projection {
804 intent: Some("test".to_string()),
805 source_facets: facets.iter().map(|s| s.to_string()).collect(),
806 reference_mems: vec![],
807 destination_mem: "v".to_string(),
808 }
809 }
810
811 fn ingest(projection: &str) -> Ingest {
812 Ingest {
813 projection: projection.to_string(),
814 mode: IngestMode::Discovery,
815 trigger: IngestTrigger::Loop,
816 batch_size: 10,
817 deny_paths: vec![],
818 }
819 }
820
821 #[test]
822 fn add_then_duplicate_medium_refuses() {
823 let tmp = TempDir::new().unwrap();
824 let root = tmp.path();
825 add_medium(root, "v", "m", &medium("m")).unwrap();
826 let err = add_medium(root, "v", "m", &medium("m")).unwrap_err();
827 assert!(
828 matches!(err, PipelineEditError::AlreadyExists { .. }),
829 "got {err:?}"
830 );
831 }
832
833 #[test]
834 fn update_missing_medium_refuses() {
835 let tmp = TempDir::new().unwrap();
836 let err = update_medium(tmp.path(), "v", "m", &medium("m")).unwrap_err();
837 assert!(
838 matches!(err, PipelineEditError::NotFound { .. }),
839 "got {err:?}"
840 );
841 }
842
843 #[test]
844 fn delete_medium_refused_while_a_facet_references_it() {
845 let tmp = TempDir::new().unwrap();
846 let root = tmp.path();
847 add_medium(root, "v", "m", &medium("m")).unwrap();
848 add_facet(root, "v", "f", &facet("f", "m")).unwrap();
849
850 let err = delete_medium(root, "v", "m").unwrap_err();
851 match err {
852 PipelineEditError::Referenced { referrers, .. } => assert_eq!(referrers, vec!["f"]),
853 other => panic!("expected Referenced, got {other:?}"),
854 }
855 delete_facet(root, "v", "f").unwrap();
857 delete_medium(root, "v", "m").unwrap();
858 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
859 assert!(configs.mediums.is_empty() && configs.facets.is_empty());
860 }
861
862 #[test]
863 fn rename_medium_repoints_dependent_facets_and_updates_embedded_name() {
864 let tmp = TempDir::new().unwrap();
865 let root = tmp.path();
866 add_medium(root, "v", "old", &medium("old")).unwrap();
867 add_facet(root, "v", "f", &facet("f", "old")).unwrap();
868
869 rename_medium(root, "v", "old", "new").unwrap();
870
871 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
872 assert_eq!(configs.mediums.len(), 1);
873 assert_eq!(configs.mediums[0].name, "new");
874 assert_eq!(configs.mediums[0].config.name, "new");
876 assert_eq!(configs.facets[0].config.medium, "new");
878 }
879
880 #[test]
881 fn rename_medium_refuses_existing_target() {
882 let tmp = TempDir::new().unwrap();
883 let root = tmp.path();
884 add_medium(root, "v", "a", &medium("a")).unwrap();
885 add_medium(root, "v", "b", &medium("b")).unwrap();
886 let err = rename_medium(root, "v", "a", "b").unwrap_err();
887 assert!(
888 matches!(err, PipelineEditError::RenameTargetExists { .. }),
889 "got {err:?}"
890 );
891 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
893 assert_eq!(configs.mediums.len(), 2);
894 }
895
896 #[test]
897 fn rename_medium_to_same_name_is_a_noop() {
898 let tmp = TempDir::new().unwrap();
899 let root = tmp.path();
900 add_medium(root, "v", "m", &medium("m")).unwrap();
901 rename_medium(root, "v", "m", "m").unwrap();
902 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
903 assert_eq!(configs.mediums.len(), 1);
904 assert_eq!(configs.mediums[0].config, medium("m"));
905 }
906
907 #[test]
908 fn delete_facet_refused_while_a_projection_references_it() {
909 let tmp = TempDir::new().unwrap();
910 let root = tmp.path();
911 add_facet(root, "v", "f", &facet("f", "m")).unwrap();
912 add_projection(root, "v", "p", &projection(&["f"])).unwrap();
913 let err = delete_facet(root, "v", "f").unwrap_err();
914 assert!(
915 matches!(err, PipelineEditError::Referenced { .. }),
916 "got {err:?}"
917 );
918 }
919
920 #[test]
921 fn rename_facet_repoints_dependent_projections() {
922 let tmp = TempDir::new().unwrap();
923 let root = tmp.path();
924 add_facet(root, "v", "old", &facet("old", "m")).unwrap();
925 add_projection(root, "v", "p", &projection(&["old", "other"])).unwrap();
926
927 rename_facet(root, "v", "old", "new").unwrap();
928
929 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
930 assert_eq!(configs.facets[0].name, "new");
931 assert_eq!(configs.facets[0].config.name, "new");
932 assert_eq!(
933 configs.projections[0].config.source_facets,
934 vec!["new", "other"]
935 );
936 }
937
938 #[test]
939 fn delete_projection_refused_while_an_ingest_runs_it() {
940 let tmp = TempDir::new().unwrap();
941 let root = tmp.path();
942 add_projection(root, "v", "p", &projection(&[])).unwrap();
943 pipeline_store::write_ingest(root, "i", &ingest("v/p")).unwrap();
944 let err = delete_projection(root, "v", "p").unwrap_err();
945 match err {
946 PipelineEditError::Referenced { referrers, .. } => assert_eq!(referrers, vec!["i"]),
947 other => panic!("expected Referenced, got {other:?}"),
948 }
949 }
950
951 #[test]
952 fn parse_json_accepts_a_valid_medium() {
953 let m: Medium =
954 parse_json(r#"{"name":"m","type":"codebase","pointer":".."}"#, "medium").unwrap();
955 assert_eq!(m.name, "m");
956 assert_eq!(m.medium_type, MediumType::Codebase);
957 }
958
959 #[test]
960 fn parse_json_maps_a_bad_payload_to_invalid_json() {
961 let err = parse_json::<Medium>("{ not json", "medium").unwrap_err();
962 assert!(
963 matches!(
964 err,
965 PipelineEditError::InvalidJson {
966 primitive: "medium",
967 ..
968 }
969 ),
970 "got {err:?}"
971 );
972 }
973
974 #[test]
975 fn rename_projection_repoints_dependent_ingests() {
976 let tmp = TempDir::new().unwrap();
977 let root = tmp.path();
978 add_projection(root, "v", "old", &projection(&[])).unwrap();
979 pipeline_store::write_ingest(root, "i", &ingest("v/old")).unwrap();
980
981 rename_projection(root, "v", "old", "new").unwrap();
982
983 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
984 assert_eq!(configs.projections[0].name, "new");
985 assert_eq!(configs.ingests[0].config.projection, "v/new");
986 }
987
988 #[test]
989 fn add_then_duplicate_ingest_refuses() {
990 let tmp = TempDir::new().unwrap();
991 let root = tmp.path();
992 add_ingest(root, "i", &ingest("v/p")).unwrap();
993 let err = add_ingest(root, "i", &ingest("v/p")).unwrap_err();
994 assert!(
995 matches!(
996 err,
997 PipelineEditError::AlreadyExists {
998 primitive: "ingest",
999 ..
1000 }
1001 ),
1002 "got {err:?}"
1003 );
1004 }
1005
1006 #[test]
1007 fn update_missing_ingest_refuses() {
1008 let tmp = TempDir::new().unwrap();
1009 let err = update_ingest(tmp.path(), "i", &ingest("v/p")).unwrap_err();
1010 assert!(
1011 matches!(
1012 err,
1013 PipelineEditError::NotFound {
1014 primitive: "ingest",
1015 ..
1016 }
1017 ),
1018 "got {err:?}"
1019 );
1020 }
1021
1022 #[test]
1023 fn update_ingest_overwrites_and_delete_removes() {
1024 let tmp = TempDir::new().unwrap();
1025 let root = tmp.path();
1026 add_ingest(root, "i", &ingest("v/p")).unwrap();
1027 let mut changed = ingest("v/p");
1028 changed.batch_size = 99;
1029 update_ingest(root, "i", &changed).unwrap();
1030 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
1031 assert_eq!(configs.ingests[0].config.batch_size, 99);
1032
1033 delete_ingest(root, "i").unwrap();
1035 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
1036 assert!(configs.ingests.is_empty());
1037 }
1038
1039 #[test]
1040 fn rename_ingest_moves_the_record() {
1041 let tmp = TempDir::new().unwrap();
1042 let root = tmp.path();
1043 add_ingest(root, "old", &ingest("v/p")).unwrap();
1044 rename_ingest(root, "old", "new").unwrap();
1045 let configs = pipeline_store::load_pipeline_configs(root).unwrap();
1046 assert_eq!(configs.ingests.len(), 1);
1047 assert_eq!(configs.ingests[0].name, "new");
1048 add_ingest(root, "other", &ingest("v/p")).unwrap();
1050 let err = rename_ingest(root, "other", "new").unwrap_err();
1051 assert!(
1052 matches!(
1053 err,
1054 PipelineEditError::RenameTargetExists {
1055 primitive: "ingest",
1056 ..
1057 }
1058 ),
1059 "got {err:?}"
1060 );
1061 }
1062}