use std::collections::{BTreeMap, BTreeSet};
use serde::Serialize;
use sha2::{Digest, Sha256};
use super::fetch::RemoteTemplate;
use super::spec::{LaunchPolicy, Origin, PrunePolicy, Sidecar};
use crate::error::{CliError, CliResult};
use crate::serve::history::templates::{TemplateId, TemplateStatus, VersionChannel};
use crate::serve::load::ConfigFormat;
#[derive(Debug, Clone, PartialEq, Eq, Serialize)]
pub struct LocalTemplate {
pub id: String,
pub status: TemplateStatus,
pub newest: Option<u32>,
pub newest_hash: Option<String>,
#[serde(default)]
pub newest_deprecated: bool,
pub stable: Option<u32>,
}
#[derive(Debug, Clone, PartialEq, Serialize)]
#[serde(tag = "action", rename_all = "snake_case")]
pub enum SyncAction {
Register {
id: String,
#[serde(skip)]
body: String,
#[serde(skip)]
format: ConfigFormat,
#[serde(skip_serializing_if = "Option::is_none")]
description: Option<String>,
launch: bool,
#[serde(skip_serializing_if = "Vec::is_empty")]
tags: Vec<VersionChannel>,
#[serde(skip_serializing_if = "Option::is_none")]
replaces: Option<u32>,
},
Launch { id: String, version: u32 },
Revive { id: String },
Unchanged { id: String, version: u32 },
Orphaned { id: String },
Deprecate { id: String },
DeprecateVersion {
id: String,
version: u32,
reason: String,
},
Skipped { name: String, reason: String },
}
impl SyncAction {
pub fn is_mutation(&self) -> bool {
matches!(
self,
Self::Register { .. }
| Self::Launch { .. }
| Self::Revive { .. }
| Self::Deprecate { .. }
| Self::DeprecateVersion { .. }
)
}
}
#[derive(Debug, Clone, PartialEq, Serialize)]
pub struct SyncPlan {
pub origin: String,
pub actions: Vec<SyncAction>,
}
impl SyncPlan {
pub fn mutations(&self) -> usize {
self.actions.iter().filter(|a| a.is_mutation()).count()
}
}
pub fn body_hash(body: &str) -> CliResult<String> {
let value: serde_yaml::Value = serde_yaml::from_str(body)
.map_err(|e| CliError::Config(format!("body is not valid YAML/JSON: {e}")))?;
let canonical = serde_yaml::to_string(&value)
.map_err(|e| CliError::Config(format!("body could not be canonicalized: {e}")))?;
let digest = Sha256::digest(canonical.as_bytes());
Ok(format!("{digest:x}"))
}
fn should_launch(policy: LaunchPolicy, sidecar: &Sidecar) -> bool {
match policy {
LaunchPolicy::Ignore => false,
LaunchPolicy::Follow => sidecar.launch,
LaunchPolicy::Always => true,
}
}
pub fn plan(origin: &Origin, remote: &[RemoteTemplate], local: &[LocalTemplate]) -> SyncPlan {
let local_by_id: BTreeMap<&str, &LocalTemplate> = local
.iter()
.filter(|l| l.id.starts_with(&origin.prefix))
.map(|l| (l.id.as_str(), l))
.collect();
let mut remote: Vec<&RemoteTemplate> = remote.iter().collect();
remote.sort_by(|a, b| a.stem.cmp(&b.stem));
let mut actions = Vec::new();
let mut seen: BTreeSet<String> = BTreeSet::new();
for r in remote {
if let Some(why) = &r.retired {
let id = format!("{}{}", origin.prefix, r.stem);
seen.insert(id.clone());
let registered = local_by_id.get(id.as_str()).and_then(|l| {
let same = l.newest_hash.is_some() && body_hash(&r.body).ok() == l.newest_hash;
(same && !l.newest_deprecated).then_some(l.newest).flatten()
});
match registered {
Some(version) => actions.push(SyncAction::DeprecateVersion {
id,
version,
reason: why.clone(),
}),
None => actions.push(SyncAction::Skipped {
name: r.stem.clone(),
reason: why.clone(),
}),
}
continue;
}
let id = format!("{}{}", origin.prefix, r.stem);
if let Err(e) = TemplateId::parse(&id) {
actions.push(SyncAction::Skipped {
name: r.stem.clone(),
reason: e.to_string(),
});
continue;
}
if !seen.insert(id.clone()) {
actions.push(SyncAction::Skipped {
name: r.stem.clone(),
reason: format!("duplicate template id '{id}'"),
});
continue;
}
let hash = match body_hash(&r.body) {
Ok(h) => h,
Err(e) => {
actions.push(SyncAction::Skipped {
name: r.stem.clone(),
reason: e.to_string(),
});
continue;
}
};
let sidecar = r.sidecar.clone().unwrap_or_default();
let tags: Vec<VersionChannel> = match sidecar
.tags
.iter()
.map(|t| VersionChannel::parse(t))
.collect::<CliResult<Vec<_>>>()
{
Ok(t) => t,
Err(e) => {
actions.push(SyncAction::Skipped {
name: r.stem.clone(),
reason: format!("sidecar tags: {e}"),
});
continue;
}
};
let launch = should_launch(origin.launch, &sidecar);
match local_by_id.get(id.as_str()) {
Some(l) if l.newest.is_some() && l.newest_hash.as_deref() == Some(hash.as_str()) => {
let newest = l.newest.expect("checked");
let mut acted = false;
if l.status == TemplateStatus::Deprecated && origin.prune == PrunePolicy::Deprecate
{
actions.push(SyncAction::Revive { id: id.clone() });
acted = true;
}
if launch && l.stable != Some(newest) {
actions.push(SyncAction::Launch {
id: id.clone(),
version: newest,
});
acted = true;
}
if !acted {
actions.push(SyncAction::Unchanged {
id,
version: newest,
});
}
}
other => actions.push(SyncAction::Register {
id,
body: r.body.clone(),
format: r.format,
description: sidecar.description.clone(),
launch,
tags,
replaces: other.and_then(|l| l.newest),
}),
}
}
for (id, l) in local_by_id {
if seen.contains(id) {
continue;
}
match origin.prune {
PrunePolicy::Keep => actions.push(SyncAction::Orphaned { id: id.to_string() }),
PrunePolicy::Deprecate if l.status == TemplateStatus::Deprecated => {
actions.push(SyncAction::Orphaned { id: id.to_string() })
}
PrunePolicy::Deprecate => actions.push(SyncAction::Deprecate { id: id.to_string() }),
}
}
SyncPlan {
origin: origin.name.clone(),
actions,
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::templates::sync::spec::{GithubSource, OriginSource};
const BODY_A: &str = "version: 1\nname: a\npipeline:\n source: {type: rest, config: {base_url: \"https://x\", path: /e}}\n sink: {type: stdout, config: {}}\n";
fn origin(prefix: &str, launch: LaunchPolicy, prune: PrunePolicy) -> Origin {
Origin {
name: "o".into(),
source: OriginSource::Github(GithubSource {
repo: "acme/t".into(),
r#ref: "main".into(),
path: String::new(),
paths: Vec::new(),
token: None,
api_base: "https://api.github.com".into(),
}),
prefix: prefix.into(),
launch,
prune,
interval_secs: None,
}
}
fn remote(stem: &str, body: &str, sidecar: Option<Sidecar>) -> RemoteTemplate {
RemoteTemplate {
stem: stem.into(),
body: body.into(),
format: ConfigFormat::Yaml,
sidecar,
retired: None,
}
}
fn local(
id: &str,
status: TemplateStatus,
newest: u32,
body: &str,
stable: Option<u32>,
) -> LocalTemplate {
LocalTemplate {
id: id.into(),
status,
newest: Some(newest),
newest_hash: Some(body_hash(body).unwrap()),
newest_deprecated: false,
stable,
}
}
fn ids(plan: &SyncPlan) -> Vec<String> {
plan.actions
.iter()
.map(|a| match a {
SyncAction::Register { id, .. }
| SyncAction::Launch { id, .. }
| SyncAction::Revive { id }
| SyncAction::Unchanged { id, .. }
| SyncAction::Orphaned { id }
| SyncAction::Deprecate { id }
| SyncAction::DeprecateVersion { id, .. } => id.clone(),
SyncAction::Skipped { name, .. } => format!("skipped:{name}"),
})
.collect()
}
#[test]
fn body_hash_is_canonical() {
let a = body_hash("version: 1\nname: a\n").unwrap();
let b = body_hash("# hello\nversion: 1\nname: a # trailing\n\n").unwrap();
let c = body_hash("version: 1\nname: b\n").unwrap();
assert_eq!(a, b);
assert_ne!(a, c);
assert_eq!(a, body_hash(r#"{"version":1,"name":"a"}"#).unwrap());
assert!(body_hash(": : :").is_err());
}
#[test]
fn new_template_registers_with_prefix_and_policy() {
let o = origin("plat-", LaunchPolicy::Ignore, PrunePolicy::Keep);
let p = plan(&o, &[remote("sync", BODY_A, None)], &[]);
assert_eq!(p.actions.len(), 1);
match &p.actions[0] {
SyncAction::Register {
id,
launch,
replaces,
tags,
..
} => {
assert_eq!(id, "plat-sync");
assert!(!launch, "launch: ignore never launches");
assert_eq!(*replaces, None);
assert!(tags.is_empty());
}
other => panic!("{other:?}"),
}
assert_eq!(p.mutations(), 1);
}
#[test]
fn unchanged_body_is_a_noop_and_changed_body_appends() {
let o = origin("", LaunchPolicy::Ignore, PrunePolicy::Keep);
let l = [local("sync", TemplateStatus::Launched, 3, BODY_A, Some(3))];
let p = plan(&o, &[remote("sync", ": : :", None)], &l);
assert!(matches!(&p.actions[0], SyncAction::Skipped { .. }), "{p:?}");
let p = plan(&o, &[remote("sync", BODY_A, None)], &l);
assert_eq!(
p.actions,
vec![SyncAction::Unchanged {
id: "sync".into(),
version: 3
}]
);
assert_eq!(p.mutations(), 0);
let changed = BODY_A.replace("path: /e", "path: /e2");
let p = plan(&o, &[remote("sync", &changed, None)], &l);
match &p.actions[0] {
SyncAction::Register { replaces, .. } => assert_eq!(*replaces, Some(3)),
other => panic!("{other:?}"),
}
}
#[test]
fn launch_policy_follow_reads_the_sidecar_and_always_ignores_it() {
let side = Sidecar {
launch: true,
description: Some("d".into()),
tags: vec!["dev".into()],
..Default::default()
};
let follow = origin("", LaunchPolicy::Follow, PrunePolicy::Keep);
let p = plan(&follow, &[remote("a", BODY_A, Some(side.clone()))], &[]);
match &p.actions[0] {
SyncAction::Register {
launch,
description,
tags,
..
} => {
assert!(launch);
assert_eq!(description.as_deref(), Some("d"));
assert_eq!(tags, &[VersionChannel::parse("dev").unwrap()]);
}
other => panic!("{other:?}"),
}
let p = plan(&follow, &[remote("a", BODY_A, None)], &[]);
assert!(matches!(
&p.actions[0],
SyncAction::Register { launch: false, .. }
));
let always = origin("", LaunchPolicy::Always, PrunePolicy::Keep);
let p = plan(&always, &[remote("a", BODY_A, None)], &[]);
assert!(matches!(
&p.actions[0],
SyncAction::Register { launch: true, .. }
));
}
#[test]
fn unchanged_but_unlaunched_plans_a_launch_under_launch_policies() {
let always = origin("", LaunchPolicy::Always, PrunePolicy::Keep);
let l = [local("a", TemplateStatus::Draft, 2, BODY_A, None)];
let p = plan(&always, &[remote("a", BODY_A, None)], &l);
assert_eq!(
p.actions,
vec![SyncAction::Launch {
id: "a".into(),
version: 2
}]
);
let l = [local("a", TemplateStatus::Launched, 2, BODY_A, Some(2))];
let p = plan(&always, &[remote("a", BODY_A, None)], &l);
assert!(matches!(&p.actions[0], SyncAction::Unchanged { .. }));
let ignore = origin("", LaunchPolicy::Ignore, PrunePolicy::Keep);
let l = [local("a", TemplateStatus::Draft, 2, BODY_A, None)];
let p = plan(&ignore, &[remote("a", BODY_A, None)], &l);
assert!(matches!(&p.actions[0], SyncAction::Unchanged { .. }));
}
#[test]
fn orphans_are_reported_under_keep_and_deprecated_under_deprecate() {
let keep = origin("p-", LaunchPolicy::Ignore, PrunePolicy::Keep);
let l = [
local("p-gone", TemplateStatus::Launched, 1, BODY_A, Some(1)),
local("other-ns", TemplateStatus::Launched, 1, BODY_A, Some(1)),
];
let p = plan(&keep, &[], &l);
assert_eq!(ids(&p), vec!["p-gone"]);
assert!(matches!(&p.actions[0], SyncAction::Orphaned { .. }));
assert_eq!(p.mutations(), 0, "keep never mutates");
let dep = origin("p-", LaunchPolicy::Ignore, PrunePolicy::Deprecate);
let p = plan(&dep, &[], &l);
assert_eq!(
p.actions,
vec![SyncAction::Deprecate {
id: "p-gone".into()
}]
);
let l = [local(
"p-gone",
TemplateStatus::Deprecated,
1,
BODY_A,
Some(1),
)];
let p = plan(&dep, &[], &l);
assert_eq!(
p.actions,
vec![SyncAction::Orphaned {
id: "p-gone".into()
}]
);
}
#[test]
fn a_returning_template_is_revived_only_when_sync_owns_deprecation() {
let l = [local("a", TemplateStatus::Deprecated, 1, BODY_A, Some(1))];
let dep = origin("", LaunchPolicy::Ignore, PrunePolicy::Deprecate);
let p = plan(&dep, &[remote("a", BODY_A, None)], &l);
assert_eq!(p.actions, vec![SyncAction::Revive { id: "a".into() }]);
let keep = origin("", LaunchPolicy::Ignore, PrunePolicy::Keep);
let p = plan(&keep, &[remote("a", BODY_A, None)], &l);
assert!(matches!(&p.actions[0], SyncAction::Unchanged { .. }));
let always = origin("", LaunchPolicy::Always, PrunePolicy::Deprecate);
let l = [local("a", TemplateStatus::Deprecated, 2, BODY_A, Some(1))];
let p = plan(&always, &[remote("a", BODY_A, None)], &l);
assert_eq!(
p.actions,
vec![
SyncAction::Revive { id: "a".into() },
SyncAction::Launch {
id: "a".into(),
version: 2
},
]
);
}
#[test]
fn a_body_the_catalog_retired_is_skipped_with_its_reason() {
let o = origin("", LaunchPolicy::Always, PrunePolicy::Keep);
let mut r = remote("acme/erp", BODY_A, None);
r.retired = Some("catalog v3 is deprecated: drops invoices".into());
let p = plan(&o, &[r], &[]);
assert_eq!(
p.actions,
vec![SyncAction::Skipped {
name: "acme/erp".into(),
reason: "catalog v3 is deprecated: drops invoices".into(),
}]
);
assert_eq!(p.mutations(), 0);
}
#[test]
fn a_retired_body_already_registered_retires_that_version() {
let o = origin("", LaunchPolicy::Always, PrunePolicy::Keep);
let mut r = remote("acme/erp", BODY_A, None);
r.retired = Some("catalog v3 is deprecated: drops invoices".into());
let l = local("acme/erp", TemplateStatus::Launched, 5, BODY_A, Some(5));
let p = plan(&o, std::slice::from_ref(&r), std::slice::from_ref(&l));
assert_eq!(
p.actions,
vec![SyncAction::DeprecateVersion {
id: "acme/erp".into(),
version: 5,
reason: "catalog v3 is deprecated: drops invoices".into(),
}]
);
assert_eq!(p.mutations(), 1);
let prune = origin("", LaunchPolicy::Always, PrunePolicy::Deprecate);
assert_eq!(
plan(&prune, std::slice::from_ref(&r), std::slice::from_ref(&l))
.actions
.len(),
1
);
let done = LocalTemplate {
newest_deprecated: true,
..l.clone()
};
assert!(matches!(
plan(&o, std::slice::from_ref(&r), &[done]).actions[0],
SyncAction::Skipped { .. }
));
let other = local(
"acme/erp",
TemplateStatus::Launched,
5,
"version: 1\nname: other\n",
Some(5),
);
assert!(matches!(
plan(&o, &[r], &[other]).actions[0],
SyncAction::Skipped { .. }
));
}
#[test]
fn bad_ids_tags_and_duplicates_are_skipped_not_fatal() {
let o = origin("", LaunchPolicy::Ignore, PrunePolicy::Keep);
let bad_tag = Sidecar {
tags: vec!["latest".into()],
..Default::default()
};
let p = plan(
&o,
&[
remote("Has Space", BODY_A, None),
remote("ok", BODY_A, Some(bad_tag)),
remote("fine", BODY_A, None),
remote("fine", BODY_A, None),
],
&[],
);
let skipped: Vec<&str> = p
.actions
.iter()
.filter_map(|a| match a {
SyncAction::Skipped { name, reason } => {
assert!(!reason.is_empty());
Some(name.as_str())
}
_ => None,
})
.collect();
assert_eq!(skipped, vec!["Has Space", "fine", "ok"]);
assert_eq!(p.mutations(), 1, "the first `fine` still registers");
}
#[test]
fn actions_serialize_with_a_tag_and_without_the_body() {
let a = SyncAction::Register {
id: "x".into(),
body: "secret-ish body".into(),
format: ConfigFormat::Yaml,
description: None,
launch: true,
tags: vec![],
replaces: Some(1),
};
let v = serde_json::to_value(&a).unwrap();
assert_eq!(v["action"], "register");
assert_eq!(v["replaces"], 1);
assert!(v.get("body").is_none());
assert!(v.get("description").is_none());
}
}