use std::path::Path;
use crate::engine::Engine;
use crate::pipeline::{Facet, Ingest, 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("{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_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_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_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_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_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_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_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_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_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_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_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_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(())
}
fn ingest_exists(c: &PipelineConfigs, name: &str) -> bool {
c.ingests.iter().any(|r| r.name == name)
}
pub fn add_ingest(root: &Path, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_pipeline_configs(root)?;
if ingest_exists(&configs, name) {
return Err(PipelineEditError::AlreadyExists {
primitive: "ingest",
key: name.to_string(),
});
}
pipeline_store::write_ingest(root, name, ingest)?;
Ok(())
}
pub fn update_ingest(root: &Path, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_pipeline_configs(root)?;
if !ingest_exists(&configs, name) {
return Err(PipelineEditError::NotFound {
primitive: "ingest",
key: name.to_string(),
});
}
pipeline_store::write_ingest(root, name, ingest)?;
Ok(())
}
pub fn delete_ingest(root: &Path, name: &str) -> Result<(), PipelineEditError> {
let configs = pipeline_store::load_pipeline_configs(root)?;
if !ingest_exists(&configs, name) {
return Err(PipelineEditError::NotFound {
primitive: "ingest",
key: name.to_string(),
});
}
pipeline_store::delete_ingest(root, name)?;
Ok(())
}
pub fn rename_ingest(root: &Path, old: &str, new: &str) -> Result<(), PipelineEditError> {
if old == new {
return Ok(());
}
let configs = pipeline_store::load_pipeline_configs(root)?;
if !ingest_exists(&configs, old) {
return Err(PipelineEditError::NotFound {
primitive: "ingest",
key: old.to_string(),
});
}
if ingest_exists(&configs, new) {
return Err(PipelineEditError::RenameTargetExists {
primitive: "ingest",
key: new.to_string(),
});
}
pipeline_store::rename_ingest(root, old, new)?;
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_pipeline_configs(root)?);
Ok(())
}
pub fn add_medium(
&mut self,
mem: &str,
name: &str,
medium: &Medium,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_medium(&root, mem, name, medium)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_medium(
&mut self,
mem: &str,
name: &str,
medium: &Medium,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_medium(&root, mem, name, medium)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_medium(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_medium(&root, mem, name)?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_medium(
&mut self,
mem: &str,
old: &str,
new: &str,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_medium(&root, mem, old, new)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_facet(
&mut self,
mem: &str,
name: &str,
facet: &Facet,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_facet(&root, mem, name, facet)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_facet(
&mut self,
mem: &str,
name: &str,
facet: &Facet,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_facet(&root, mem, name, facet)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_facet(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_facet(&root, mem, name)?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_facet(
&mut self,
mem: &str,
old: &str,
new: &str,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_facet(&root, mem, old, new)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_projection(
&mut self,
mem: &str,
name: &str,
projection: &Projection,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_projection(&root, mem, name, projection)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_projection(
&mut self,
mem: &str,
name: &str,
projection: &Projection,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_projection(&root, mem, name, projection)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_projection(&mut self, mem: &str, name: &str) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_projection(&root, mem, name)?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_projection(
&mut self,
mem: &str,
old: &str,
new: &str,
) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_projection(&root, mem, old, new)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_ingest(&mut self, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
add_ingest(&root, name, ingest)?;
self.refresh_pipeline_configs(&root)
}
pub fn update_ingest(&mut self, name: &str, ingest: &Ingest) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
update_ingest(&root, name, ingest)?;
self.refresh_pipeline_configs(&root)
}
pub fn delete_ingest(&mut self, name: &str) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
delete_ingest(&root, name)?;
self.refresh_pipeline_configs(&root)
}
pub fn rename_ingest(&mut self, old: &str, new: &str) -> Result<(), PipelineEditError> {
let root = self.pipeline_edit_root()?;
rename_ingest(&root, old, new)?;
self.refresh_pipeline_configs(&root)
}
pub fn add_medium_json(
&mut self,
mem: &str,
name: &str,
medium_json: &str,
) -> Result<(), PipelineEditError> {
self.add_medium(mem, name, &parse_json(medium_json, "medium")?)
}
pub fn update_medium_json(
&mut self,
mem: &str,
name: &str,
medium_json: &str,
) -> Result<(), PipelineEditError> {
self.update_medium(mem, name, &parse_json(medium_json, "medium")?)
}
pub fn add_facet_json(
&mut self,
mem: &str,
name: &str,
facet_json: &str,
) -> Result<(), PipelineEditError> {
self.add_facet(mem, name, &parse_json(facet_json, "facet")?)
}
pub fn update_facet_json(
&mut self,
mem: &str,
name: &str,
facet_json: &str,
) -> Result<(), PipelineEditError> {
self.update_facet(mem, name, &parse_json(facet_json, "facet")?)
}
pub fn add_projection_json(
&mut self,
mem: &str,
name: &str,
projection_json: &str,
) -> Result<(), PipelineEditError> {
self.add_projection(mem, name, &parse_json(projection_json, "projection")?)
}
pub fn update_projection_json(
&mut self,
mem: &str,
name: &str,
projection_json: &str,
) -> Result<(), PipelineEditError> {
self.update_projection(mem, name, &parse_json(projection_json, "projection")?)
}
pub fn add_ingest_json(
&mut self,
name: &str,
ingest_json: &str,
) -> Result<(), PipelineEditError> {
self.add_ingest(name, &parse_json(ingest_json, "ingest")?)
}
pub fn update_ingest_json(
&mut self,
name: &str,
ingest_json: &str,
) -> Result<(), PipelineEditError> {
self.update_ingest(name, &parse_json(ingest_json, "ingest")?)
}
}
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::{
Ingest, IngestMode, IngestTrigger, MediumType, PatternEntry, PatternMode,
};
use tempfile::TempDir;
fn medium(name: &str) -> Medium {
Medium {
name: name.to_string(),
medium_type: MediumType::Codebase,
pointer: "../src".to_string(),
}
}
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(),
}
}
fn ingest(projection: &str) -> Ingest {
Ingest {
projection: projection.to_string(),
mode: IngestMode::Discovery,
trigger: IngestTrigger::Loop,
batch_size: 10,
deny_paths: vec![],
}
}
#[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_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_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_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_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_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_pipeline_configs(root).unwrap();
assert_eq!(configs.projections[0].name, "new");
assert_eq!(configs.ingests[0].config.projection, "v/new");
}
#[test]
fn add_then_duplicate_ingest_refuses() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_ingest(root, "i", &ingest("v/p")).unwrap();
let err = add_ingest(root, "i", &ingest("v/p")).unwrap_err();
assert!(
matches!(
err,
PipelineEditError::AlreadyExists {
primitive: "ingest",
..
}
),
"got {err:?}"
);
}
#[test]
fn update_missing_ingest_refuses() {
let tmp = TempDir::new().unwrap();
let err = update_ingest(tmp.path(), "i", &ingest("v/p")).unwrap_err();
assert!(
matches!(
err,
PipelineEditError::NotFound {
primitive: "ingest",
..
}
),
"got {err:?}"
);
}
#[test]
fn update_ingest_overwrites_and_delete_removes() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_ingest(root, "i", &ingest("v/p")).unwrap();
let mut changed = ingest("v/p");
changed.batch_size = 99;
update_ingest(root, "i", &changed).unwrap();
let configs = pipeline_store::load_pipeline_configs(root).unwrap();
assert_eq!(configs.ingests[0].config.batch_size, 99);
delete_ingest(root, "i").unwrap();
let configs = pipeline_store::load_pipeline_configs(root).unwrap();
assert!(configs.ingests.is_empty());
}
#[test]
fn rename_ingest_moves_the_record() {
let tmp = TempDir::new().unwrap();
let root = tmp.path();
add_ingest(root, "old", &ingest("v/p")).unwrap();
rename_ingest(root, "old", "new").unwrap();
let configs = pipeline_store::load_pipeline_configs(root).unwrap();
assert_eq!(configs.ingests.len(), 1);
assert_eq!(configs.ingests[0].name, "new");
add_ingest(root, "other", &ingest("v/p")).unwrap();
let err = rename_ingest(root, "other", "new").unwrap_err();
assert!(
matches!(
err,
PipelineEditError::RenameTargetExists {
primitive: "ingest",
..
}
),
"got {err:?}"
);
}
}