Skip to main content

memstead_base/
pipeline_edit.rs

1//! Edit operations over the four-primitive pipeline store, with referential
2//! integrity.
3//!
4//! Sits above the dumb file ops in [`crate::pipeline_store`]: each function
5//! loads the current store, enforces the Medium ← Facet ← Projection ← Ingest
6//! reference model (no clobber, no dangling reference), then writes. The
7//! engine's in-memory `pipeline_configs` cache is refreshed by the `Engine`
8//! wrapper methods that call these — not here — so these functions stay pure
9//! disk ops that unit-test with a bare `TempDir`.
10//!
11//! Identity is the file stem `(mem, name)`. Mediums and facets additionally
12//! carry an embedded `name` field kept equal to the stem (facets reference
13//! mediums by name, projections reference facets by name); rename here updates
14//! the file location, the embedded field, and every dependent reference
15//! together, so the store never holds a self-inconsistent or dangling record.
16//!
17//! Scope: medium / facet / projection — the editing chain the macOS app's
18//! pipeline editor drives. Ingests are read for integrity (a projection can't
19//! be deleted while an ingest runs it) and repointed on a projection rename,
20//! but interactive ingest editing has no consumer yet and is not exposed here.
21
22use 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/// Failure modes of a pipeline edit. Distinct from the entity-centric
30/// [`crate::engine::EngineError`] — these describe four-primitive store edits.
31/// `key` is the display identity: `"<mem>/<name>"`.
32#[derive(Debug, thiserror::Error)]
33pub enum PipelineEditError {
34    /// The engine was not booted from a workspace root, so there is no
35    /// four-primitive store to edit (e.g. an engine built from a bare mount
36    /// list in a test or in-memory consumer).
37    #[error("engine has no workspace root — pipeline edits require a workspace-backed engine")]
38    NoWorkspaceRoot,
39    /// A create targeted a `(mem, name)` that already holds a record.
40    #[error("{primitive} '{key}' already exists")]
41    AlreadyExists {
42        primitive: &'static str,
43        key: String,
44    },
45    /// An update / delete / rename targeted a record that does not exist.
46    #[error("{primitive} '{key}' does not exist")]
47    NotFound {
48        primitive: &'static str,
49        key: String,
50    },
51    /// A delete was refused because other records still reference the target.
52    #[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    /// A rename target `(mem, new)` already holds a record.
59    #[error("rename target {primitive} '{key}' already exists")]
60    RenameTargetExists {
61        primitive: &'static str,
62        key: String,
63    },
64    /// A JSON-string edit entry point received a payload that did not
65    /// deserialize into the target primitive.
66    #[error("invalid {primitive} JSON: {message}")]
67    InvalidJson {
68        primitive: &'static str,
69        message: String,
70    },
71    /// Underlying store IO / parse failure.
72    #[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
92/// Facet names (same mem) whose `medium` points at `name`.
93fn 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
101/// Projection names (same mem) whose `source_facets` contain `name`.
102fn 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
110/// Ingest names whose `projection` points at `<mem>/<name>`.
111fn 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
120// --- Medium ----------------------------------------------------------------
121
122/// Create a medium. Refuses if `(mem, name)` already holds one.
123pub 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
140/// Overwrite an existing medium. Refuses if `(mem, name)` does not exist.
141pub 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
158/// Delete a medium. Refuses if any facet in the same mem still references it.
159pub 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
179/// Rename a medium within its mem, updating its embedded `name` and every
180/// dependent facet's `medium` reference. No-op when `old == new`.
181pub 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    // Write the new stem first, then repoint referrers, then drop the old —
208    // no point in the sequence leaves a facet pointing at a missing medium.
209    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
225// --- Facet -----------------------------------------------------------------
226
227/// Create a facet. Refuses if `(mem, name)` already holds one.
228pub 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
245/// Overwrite an existing facet. Refuses if `(mem, name)` does not exist.
246pub 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
263/// Delete a facet. Refuses if any projection in the same mem references it.
264pub 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
284/// Rename a facet within its mem, updating its embedded `name` and every
285/// dependent projection's `source_facets` entry. No-op when `old == new`.
286pub 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
327// --- Projection ------------------------------------------------------------
328
329/// Create a projection. Refuses if `(mem, name)` already holds one.
330pub 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
347/// Overwrite an existing projection. Refuses if `(mem, name)` does not exist.
348pub 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
365/// Delete a projection. Refuses if any ingest still runs it.
366pub 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
386/// Rename a projection within its mem (a projection has no embedded name, so
387/// the file moves), repointing every ingest whose `projection` was
388/// `<mem>/<old>`. No-op when `old == new`.
389pub 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
423// --- Ingest ----------------------------------------------------------------
424//
425// Ingests are *flat* (workspace-level, not mem-scoped) and are the leaf of
426// the pipeline — nothing references an ingest — so add/update/delete only
427// check the ingest's own existence; there is no referrer gate on delete. An
428// ingest's `projection` field points at a `<mem>/<name>` projection key
429// (delete_projection enforces the inverse integrity).
430
431fn ingest_exists(c: &PipelineConfigs, name: &str) -> bool {
432    c.ingests.iter().any(|r| r.name == name)
433}
434
435/// Create an ingest. Refuses if `name` already holds one.
436pub 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
448/// Overwrite an existing ingest. Refuses if `name` does not exist.
449pub 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
461/// Delete an ingest. Nothing references an ingest, so no referrer gate.
462pub 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
474/// Rename an ingest (file-stem identity; nothing depends on it). No-op when
475/// `old == new`.
476pub 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
497// --- Engine surface --------------------------------------------------------
498//
499// Thin wrappers that route an edit through the free functions above (disk +
500// referential integrity) and then refresh the engine's in-memory
501// `pipeline_configs` snapshot so a subsequent `pipeline_configs()` read sees
502// the change. They use only the engine's public accessors, so this block stays
503// out of the `engine` module's internals.
504
505impl 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    /// Create a medium and refresh the in-memory snapshot. See [`add_medium`].
518    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    /// Overwrite a medium and refresh the snapshot. See [`update_medium`].
530    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    /// Delete a medium and refresh the snapshot. See [`delete_medium`].
542    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    /// Rename a medium and refresh the snapshot. See [`rename_medium`].
549    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    /// Create a facet and refresh the snapshot. See [`add_facet`].
561    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    /// Overwrite a facet and refresh the snapshot. See [`update_facet`].
573    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    /// Delete a facet and refresh the snapshot. See [`delete_facet`].
585    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    /// Rename a facet and refresh the snapshot. See [`rename_facet`].
592    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    /// Create a projection and refresh the snapshot. See [`add_projection`].
604    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    /// Overwrite a projection and refresh the snapshot. See [`update_projection`].
616    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    /// Delete a projection and refresh the snapshot. See [`delete_projection`].
628    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    /// Rename a projection and refresh the snapshot. See [`rename_projection`].
635    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    /// Create an ingest and refresh the snapshot. See [`add_ingest`].
647    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    /// Overwrite an ingest and refresh the snapshot. See [`update_ingest`].
654    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    /// Delete an ingest and refresh the snapshot. See [`delete_ingest`].
661    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    /// Rename an ingest and refresh the snapshot. See [`rename_ingest`].
668    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    // JSON-string entry points for serialization-boundary callers (UniFFI,
675    // CLI) that carry a primitive as JSON rather than a typed value. They
676    // deserialize here — where serde already lives — and delegate to the
677    // typed methods above, so the FFI translation layer needs no JSON
678    // dependency of its own. Only `add`/`update` carry a payload; `delete`
679    // and `rename` take plain string identifiers and use the typed methods
680    // directly.
681
682    /// [`Self::add_medium`] from a JSON-encoded [`Medium`].
683    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    /// [`Self::update_medium`] from a JSON-encoded [`Medium`].
693    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    /// [`Self::add_facet`] from a JSON-encoded [`Facet`].
703    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    /// [`Self::update_facet`] from a JSON-encoded [`Facet`].
713    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    /// [`Self::add_projection`] from a JSON-encoded [`Projection`].
723    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    /// [`Self::update_projection`] from a JSON-encoded [`Projection`].
733    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    /// [`Self::add_ingest`] from a JSON-encoded [`Ingest`].
743    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    /// [`Self::update_ingest`] from a JSON-encoded [`Ingest`].
752    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
761/// Deserialize a pipeline primitive from JSON, mapping a parse failure to a
762/// typed [`PipelineEditError::InvalidJson`] naming the primitive.
763fn 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        // Removing the facet frees the medium.
856        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        // Embedded name tracks the stem.
875        assert_eq!(configs.mediums[0].config.name, "new");
876        // Dependent facet now points at the new medium name.
877        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        // Nothing lost.
892        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        // Nothing references an ingest — delete needs no referrer gate.
1034        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        // Existing target refuses.
1049        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}