use std::path::Path;
use crate::engine::Engine;
use crate::pipeline::{Facet, Medium, Projection};
use crate::pipeline_store::{self, PipelineConfigs};
use crate::workspace_store::StoreError;
#[derive(Debug, thiserror::Error)]
pub enum PipelineEditError {
#[error("engine has no workspace root — pipeline edits require a workspace-backed engine")]
NoWorkspaceRoot,
#[error("pipeline edit landed, but recording provenance failed: {0}")]
Provenance(String),
#[error("{primitive} '{key}' already exists")]
AlreadyExists {
primitive: &'static str,
key: String,
},
#[error("{primitive} '{key}' does not exist")]
NotFound {
primitive: &'static str,
key: String,
},
#[error("{primitive} '{key}' is referenced by {referrers:?} — remove or repoint them first")]
Referenced {
primitive: &'static str,
key: String,
referrers: Vec<String>,
},
#[error("rename target {primitive} '{key}' already exists")]
RenameTargetExists {
primitive: &'static str,
key: String,
},
#[error("invalid {primitive} JSON: {message}")]
InvalidJson {
primitive: &'static str,
message: String,
},
#[error(transparent)]
Store(#[from] StoreError),
}
fn key(mem: &str, name: &str) -> String {
format!("{mem}/{name}")
}
fn medium_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
c.mediums.iter().any(|r| r.mem == mem && r.name == name)
}
fn facet_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
c.facets.iter().any(|r| r.mem == mem && r.name == name)
}
fn projection_exists(c: &PipelineConfigs, mem: &str, name: &str) -> bool {
c.projections.iter().any(|r| r.mem == mem && r.name == name)
}
fn facets_referencing_medium(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
c.facets
.iter()
.filter(|r| r.mem == mem && r.config.medium == name)
.map(|r| r.name.clone())
.collect()
}
fn projections_referencing_facet(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
c.projections
.iter()
.filter(|r| r.mem == mem && r.config.source_facets.iter().any(|f| f == name))
.map(|r| r.name.clone())
.collect()
}
fn ingests_referencing_projection(c: &PipelineConfigs, mem: &str, name: &str) -> Vec<String> {
let target = key(mem, name);
c.ingests
.iter()
.filter(|r| r.config.projection == target)
.map(|r| r.name.clone())
.collect()
}
pub fn add_medium(
root: &Path,
mem: &str,
name: &str,
medium: &Medium,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if medium_exists(&configs, mem, name) {
return Err(PipelineEditError::AlreadyExists {
primitive: "medium",
key: key(mem, name),
});
}
pipeline_store::write_medium(root, mem, name, medium)?;
Ok(())
}
pub fn update_medium(
root: &Path,
mem: &str,
name: &str,
medium: &Medium,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !medium_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "medium",
key: key(mem, name),
});
}
pipeline_store::write_medium(root, mem, name, medium)?;
Ok(())
}
pub fn delete_medium(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !medium_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "medium",
key: key(mem, name),
});
}
let referrers = facets_referencing_medium(&configs, mem, name);
if !referrers.is_empty() {
return Err(PipelineEditError::Referenced {
primitive: "medium",
key: key(mem, name),
referrers,
});
}
pipeline_store::delete_medium(root, mem, name)?;
Ok(())
}
pub fn rename_medium(
root: &Path,
mem: &str,
old: &str,
new: &str,
) -> Result<(), PipelineEditError> {
if old == new {
return Ok(());
}
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
let existing = configs
.mediums
.iter()
.find(|r| r.mem == mem && r.name == old)
.ok_or_else(|| PipelineEditError::NotFound {
primitive: "medium",
key: key(mem, old),
})?;
if medium_exists(&configs, mem, new) {
return Err(PipelineEditError::RenameTargetExists {
primitive: "medium",
key: key(mem, new),
});
}
let mut renamed = existing.config.clone();
renamed.name = new.to_string();
pipeline_store::write_medium(root, mem, new, &renamed)?;
for facet in facets_referencing_medium(&configs, mem, old) {
if let Some(rec) = configs
.facets
.iter()
.find(|r| r.mem == mem && r.name == facet)
{
let mut updated = rec.config.clone();
updated.medium = new.to_string();
pipeline_store::write_facet(root, mem, &facet, &updated)?;
}
}
pipeline_store::delete_medium(root, mem, old)?;
Ok(())
}
pub fn add_facet(
root: &Path,
mem: &str,
name: &str,
facet: &Facet,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if facet_exists(&configs, mem, name) {
return Err(PipelineEditError::AlreadyExists {
primitive: "facet",
key: key(mem, name),
});
}
pipeline_store::write_facet(root, mem, name, facet)?;
Ok(())
}
pub fn update_facet(
root: &Path,
mem: &str,
name: &str,
facet: &Facet,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !facet_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "facet",
key: key(mem, name),
});
}
pipeline_store::write_facet(root, mem, name, facet)?;
Ok(())
}
pub fn delete_facet(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !facet_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "facet",
key: key(mem, name),
});
}
let referrers = projections_referencing_facet(&configs, mem, name);
if !referrers.is_empty() {
return Err(PipelineEditError::Referenced {
primitive: "facet",
key: key(mem, name),
referrers,
});
}
pipeline_store::delete_facet(root, mem, name)?;
Ok(())
}
pub fn rename_facet(root: &Path, mem: &str, old: &str, new: &str) -> Result<(), PipelineEditError> {
if old == new {
return Ok(());
}
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
let existing = configs
.facets
.iter()
.find(|r| r.mem == mem && r.name == old)
.ok_or_else(|| PipelineEditError::NotFound {
primitive: "facet",
key: key(mem, old),
})?;
if facet_exists(&configs, mem, new) {
return Err(PipelineEditError::RenameTargetExists {
primitive: "facet",
key: key(mem, new),
});
}
let mut renamed = existing.config.clone();
renamed.name = new.to_string();
pipeline_store::write_facet(root, mem, new, &renamed)?;
for proj in projections_referencing_facet(&configs, mem, old) {
if let Some(rec) = configs
.projections
.iter()
.find(|r| r.mem == mem && r.name == proj)
{
let mut updated = rec.config.clone();
for f in updated.source_facets.iter_mut() {
if f == old {
*f = new.to_string();
}
}
pipeline_store::write_projection(root, mem, &proj, &updated)?;
}
}
pipeline_store::delete_facet(root, mem, old)?;
Ok(())
}
pub fn add_projection(
root: &Path,
mem: &str,
name: &str,
projection: &Projection,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if projection_exists(&configs, mem, name) {
return Err(PipelineEditError::AlreadyExists {
primitive: "projection",
key: key(mem, name),
});
}
pipeline_store::write_projection(root, mem, name, projection)?;
Ok(())
}
pub fn update_projection(
root: &Path,
mem: &str,
name: &str,
projection: &Projection,
) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !projection_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "projection",
key: key(mem, name),
});
}
pipeline_store::write_projection(root, mem, name, projection)?;
Ok(())
}
pub fn delete_projection(root: &Path, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !projection_exists(&configs, mem, name) {
return Err(PipelineEditError::NotFound {
primitive: "projection",
key: key(mem, name),
});
}
let referrers = ingests_referencing_projection(&configs, mem, name);
if !referrers.is_empty() {
return Err(PipelineEditError::Referenced {
primitive: "projection",
key: key(mem, name),
referrers,
});
}
pipeline_store::delete_projection(root, mem, name)?;
Ok(())
}
pub fn rename_projection(
root: &Path,
mem: &str,
old: &str,
new: &str,
) -> Result<(), PipelineEditError> {
if old == new {
return Ok(());
}
let configs = pipeline_store::load_legacy_pipeline_configs(root)?;
if !projection_exists(&configs, mem, old) {
return Err(PipelineEditError::NotFound {
primitive: "projection",
key: key(mem, old),
});
}
if projection_exists(&configs, mem, new) {
return Err(PipelineEditError::RenameTargetExists {
primitive: "projection",
key: key(mem, new),
});
}
pipeline_store::rename_projection(root, mem, old, new)?;
let new_ref = key(mem, new);
for ingest in ingests_referencing_projection(&configs, mem, old) {
if let Some(rec) = configs.ingests.iter().find(|r| r.name == ingest) {
let mut updated = rec.config.clone();
updated.projection = new_ref.clone();
pipeline_store::write_ingest(root, &ingest, &updated)?;
}
}
Ok(())
}
impl Engine {
fn pipeline_edit_root(&self) -> Result<std::path::PathBuf, PipelineEditError> {
self.workspace_root()
.map(Path::to_path_buf)
.ok_or(PipelineEditError::NoWorkspaceRoot)
}
fn refresh_pipeline_configs(&mut self, root: &Path) -> Result<(), PipelineEditError> {
self.set_pipeline_configs(pipeline_store::load_legacy_pipeline_configs(root)?);
Ok(())
}
fn pipeline_provenance(
&self,
mem: &str,
kind: &str,
edits: &[(String, Option<Vec<u8>>)],
note: Option<&str>,
verb: &str,
) -> Result<(), PipelineEditError> {
self.record_pipeline_edit_provenance(mem, kind, edits, note, verb)
.map_err(|e| PipelineEditError::Provenance(e.to_string()))
}
pub fn add_medium(
&mut self,
mem: &str,
name: &str,
medium: &Medium,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_medium(&root, mem, name, medium)?;
let bytes =
serde_json::to_vec_pretty(medium).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"mediums",
&[(name.to_string(), Some(bytes))],
note,
"add",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_medium(
&mut self,
mem: &str,
name: &str,
medium: &Medium,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_medium(&root, mem, name, medium)?;
let bytes =
serde_json::to_vec_pretty(medium).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"mediums",
&[(name.to_string(), Some(bytes))],
note,
"update",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_medium(
&mut self,
mem: &str,
name: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_medium(&root, mem, name)?;
self.pipeline_provenance(mem, "mediums", &[(name.to_string(), None)], note, "delete")?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_medium(
&mut self,
mem: &str,
old: &str,
new: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_medium(&root, mem, old, new)?;
let bytes = self
.pipeline_configs()
.mediums
.iter()
.find(|r| r.mem == mem && r.name == old)
.map(|r| serde_json::to_vec_pretty(&r.config))
.transpose()
.map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"mediums",
&[(old.to_string(), None), (new.to_string(), bytes)],
note,
"rename",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_facet(
&mut self,
mem: &str,
name: &str,
facet: &Facet,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_facet(&root, mem, name, facet)?;
let bytes =
serde_json::to_vec_pretty(facet).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"facets",
&[(name.to_string(), Some(bytes))],
note,
"add",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_facet(
&mut self,
mem: &str,
name: &str,
facet: &Facet,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_facet(&root, mem, name, facet)?;
let bytes =
serde_json::to_vec_pretty(facet).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"facets",
&[(name.to_string(), Some(bytes))],
note,
"update",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_facet(
&mut self,
mem: &str,
name: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_facet(&root, mem, name)?;
self.pipeline_provenance(mem, "facets", &[(name.to_string(), None)], note, "delete")?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_facet(
&mut self,
mem: &str,
old: &str,
new: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_facet(&root, mem, old, new)?;
let bytes = self
.pipeline_configs()
.facets
.iter()
.find(|r| r.mem == mem && r.name == old)
.map(|r| serde_json::to_vec_pretty(&r.config))
.transpose()
.map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"facets",
&[(old.to_string(), None), (new.to_string(), bytes)],
note,
"rename",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_projection(
&mut self,
mem: &str,
name: &str,
projection: &Projection,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_projection(&root, mem, name, projection)?;
let bytes =
serde_json::to_vec_pretty(projection).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"projections",
&[(name.to_string(), Some(bytes))],
note,
"add",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_projection(
&mut self,
mem: &str,
name: &str,
projection: &Projection,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_projection(&root, mem, name, projection)?;
let bytes =
serde_json::to_vec_pretty(projection).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"projections",
&[(name.to_string(), Some(bytes))],
note,
"update",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_projection(
&mut self,
mem: &str,
name: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_projection(&root, mem, name)?;
self.pipeline_provenance(
mem,
"projections",
&[(name.to_string(), None)],
note,
"delete",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_projection(
&mut self,
mem: &str,
old: &str,
new: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_projection(&root, mem, old, new)?;
let bytes = self
.pipeline_configs()
.projections
.iter()
.find(|r| r.mem == mem && r.name == old)
.map(|r| serde_json::to_vec_pretty(&r.config))
.transpose()
.map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"projections",
&[(old.to_string(), None), (new.to_string(), bytes)],
note,
"rename",
)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_medium_json(
&mut self,
mem: &str,
name: &str,
medium_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
self.add_medium(mem, name, &parse_json(medium_json, "medium")?, note)
}
pub fn update_medium_json(
&mut self,
mem: &str,
name: &str,
medium_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
self.update_medium(mem, name, &parse_json(medium_json, "medium")?, note)
}
pub fn add_facet_json(
&mut self,
mem: &str,
name: &str,
facet_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
self.add_facet(mem, name, &parse_json(facet_json, "facet")?, note)
}
pub fn update_facet_json(
&mut self,
mem: &str,
name: &str,
facet_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
self.update_facet(mem, name, &parse_json(facet_json, "facet")?, note)
}
pub fn add_projection_json(
&mut self,
mem: &str,
name: &str,
projection_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
let incoming: Projection = parse_json(projection_json, "projection")?;
let binding = crate::binding::BindingV1 {
version: crate::binding::BINDING_VERSION,
intent: incoming.intent,
source_facets: incoming.source_facets,
reference_mems: incoming.reference_mems,
destination_mem: incoming.destination_mem,
deny_paths: Vec::new(),
coverage_semantics: crate::binding::CoverageSemantics::default(),
rules: incoming.rules,
prune: None,
operations: crate::binding::Operations {
build: Some(crate::binding::BuildOperation {
mode: crate::binding::BuildMode::Discovery,
trigger: crate::pipeline::IngestTrigger::Loop,
batch_size: 20,
post_actions: None,
}),
sync: None,
verify: None,
},
};
self.write_binding_edit(mem, name, &binding, &root, note, "add")
}
pub fn update_projection_json(
&mut self,
mem: &str,
name: &str,
projection_json: &str,
note: Option<&str>,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
let incoming: Projection = parse_json(projection_json, "projection")?;
let mut binding = pipeline_store::read_binding(&root, mem, name)?;
binding.intent = incoming.intent;
binding.source_facets = incoming.source_facets;
binding.reference_mems = incoming.reference_mems;
binding.destination_mem = incoming.destination_mem;
if incoming.rules.is_some() {
binding.rules = incoming.rules;
}
self.write_binding_edit(mem, name, &binding, &root, note, "update")
}
fn write_binding_edit(
&mut self,
mem: &str,
name: &str,
binding: &crate::binding::BindingV1,
root: &Path,
note: Option<&str>,
verb: &str,
) -> Result<(), PipelineEditError> {
pipeline_store::write_binding(root, mem, name, binding)?;
let bytes =
serde_json::to_vec_pretty(binding).map_err(|e| PipelineEditError::InvalidJson {
primitive: "config",
message: e.to_string(),
})?;
self.pipeline_provenance(
mem,
"projections",
&[(name.to_string(), Some(bytes))],
note,
verb,
)?;
self.refresh_pipeline_configs(root)
}
}
fn parse_json<T: serde::de::DeserializeOwned>(
json: &str,
primitive: &'static str,
) -> Result<T, PipelineEditError> {
serde_json::from_str(json).map_err(|e| PipelineEditError::InvalidJson {
primitive,
message: e.to_string(),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::pipeline::{IngestTrigger, MediumType, PatternEntry, PatternMode};
use crate::pipeline_store::{LegacyIngest, LegacyIngestMode};
use tempfile::TempDir;
fn medium(name: &str) -> Medium {
Medium {
name: name.to_string(),
medium_type: MediumType::Codebase,
pointer: "../src".to_string(),
change_detection: None,
}
}
fn facet(name: &str, medium: &str) -> Facet {
Facet {
name: name.to_string(),
medium: medium.to_string(),
scope: vec![PatternEntry {
path: "**/*.rs".to_string(),
mode: PatternMode::Allow,
}],
engagement: None,
preparation: None,
}
}
fn projection(facets: &[&str]) -> Projection {
Projection {
intent: Some("test".to_string()),
source_facets: facets.iter().map(|s| s.to_string()).collect(),
reference_mems: vec![],
destination_mem: "v".to_string(),
rules: None,
}
}
fn ingest(projection: &str) -> LegacyIngest {
LegacyIngest {
projection: projection.to_string(),
mode: LegacyIngestMode::Discovery,
trigger: IngestTrigger::Loop,
batch_size: 10,
deny_paths: vec![],
post_actions: None,
}
}
#[test]
fn add_then_duplicate_medium_refuses() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_medium(root, "v", "m", &medium("m")).unwrap();
let err = add_medium(root, "v", "m", &medium("m")).unwrap_err();
assert!(
matches!(err, PipelineEditError::AlreadyExists { .. }),
"got {err:?}"
);
}
#[test]
fn update_missing_medium_refuses() {
let tmp = TempDir::new().unwrap();
let err = update_medium(tmp.path(), "v", "m", &medium("m")).unwrap_err();
assert!(
matches!(err, PipelineEditError::NotFound { .. }),
"got {err:?}"
);
}
#[test]
fn delete_medium_refused_while_a_facet_references_it() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_medium(root, "v", "m", &medium("m")).unwrap();
add_facet(root, "v", "f", &facet("f", "m")).unwrap();
let err = delete_medium(root, "v", "m").unwrap_err();
match err {
PipelineEditError::Referenced { referrers, .. } => assert_eq!(referrers, vec!["f"]),
other => panic!("expected Referenced, got {other:?}"),
}
delete_facet(root, "v", "f").unwrap();
delete_medium(root, "v", "m").unwrap();
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert!(configs.mediums.is_empty() && configs.facets.is_empty());
}
#[test]
fn rename_medium_repoints_dependent_facets_and_updates_embedded_name() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_medium(root, "v", "old", &medium("old")).unwrap();
add_facet(root, "v", "f", &facet("f", "old")).unwrap();
rename_medium(root, "v", "old", "new").unwrap();
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert_eq!(configs.mediums.len(), 1);
assert_eq!(configs.mediums[0].name, "new");
assert_eq!(configs.mediums[0].config.name, "new");
assert_eq!(configs.facets[0].config.medium, "new");
}
#[test]
fn rename_medium_refuses_existing_target() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_medium(root, "v", "a", &medium("a")).unwrap();
add_medium(root, "v", "b", &medium("b")).unwrap();
let err = rename_medium(root, "v", "a", "b").unwrap_err();
assert!(
matches!(err, PipelineEditError::RenameTargetExists { .. }),
"got {err:?}"
);
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert_eq!(configs.mediums.len(), 2);
}
#[test]
fn rename_medium_to_same_name_is_a_noop() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_medium(root, "v", "m", &medium("m")).unwrap();
rename_medium(root, "v", "m", "m").unwrap();
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert_eq!(configs.mediums.len(), 1);
assert_eq!(configs.mediums[0].config, medium("m"));
}
#[test]
fn delete_facet_refused_while_a_projection_references_it() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_facet(root, "v", "f", &facet("f", "m")).unwrap();
add_projection(root, "v", "p", &projection(&["f"])).unwrap();
let err = delete_facet(root, "v", "f").unwrap_err();
assert!(
matches!(err, PipelineEditError::Referenced { .. }),
"got {err:?}"
);
}
#[test]
fn rename_facet_repoints_dependent_projections() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_facet(root, "v", "old", &facet("old", "m")).unwrap();
add_projection(root, "v", "p", &projection(&["old", "other"])).unwrap();
rename_facet(root, "v", "old", "new").unwrap();
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert_eq!(configs.facets[0].name, "new");
assert_eq!(configs.facets[0].config.name, "new");
assert_eq!(
configs.projections[0].config.source_facets,
vec!["new", "other"]
);
}
#[test]
fn delete_projection_refused_while_an_ingest_runs_it() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_projection(root, "v", "p", &projection(&[])).unwrap();
pipeline_store::write_ingest(root, "i", &ingest("v/p")).unwrap();
let err = delete_projection(root, "v", "p").unwrap_err();
match err {
PipelineEditError::Referenced { referrers, .. } => assert_eq!(referrers, vec!["i"]),
other => panic!("expected Referenced, got {other:?}"),
}
}
#[test]
fn parse_json_accepts_a_valid_medium() {
let m: Medium =
parse_json(r#"{"name":"m","type":"codebase","pointer":".."}"#, "medium").unwrap();
assert_eq!(m.name, "m");
assert_eq!(m.medium_type, MediumType::Codebase);
}
#[test]
fn parse_json_maps_a_bad_payload_to_invalid_json() {
let err = parse_json::<Medium>("{ not json", "medium").unwrap_err();
assert!(
matches!(
err,
PipelineEditError::InvalidJson {
primitive: "medium",
..
}
),
"got {err:?}"
);
}
#[test]
fn rename_projection_repoints_dependent_ingests() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_projection(root, "v", "old", &projection(&[])).unwrap();
pipeline_store::write_ingest(root, "i", &ingest("v/old")).unwrap();
rename_projection(root, "v", "old", "new").unwrap();
let configs = pipeline_store::load_legacy_pipeline_configs(root).unwrap();
assert_eq!(configs.projections[0].name, "new");
assert_eq!(configs.ingests[0].config.projection, "v/new");
}
}