From d1ad08ef3a324ea5df4c555dc51527367780df25 Mon Sep 17 00:00:00 2001 From: Catty Steve <4795515+Catty2014@user.noreply.gitee.com> Date: Thu, 21 May 2026 20:04:47 +0800 Subject: [PATCH] feat: add workshop-baker-params for config layering Implement a priority-based merge system for pipeline configuration. Parameters carry their source (Base, Courier, Bakerd) and only higher-priority writers override. Migrate notification, prebake, bake, and finalize config to use ParamVal wrappers. Remove hardcoded constants and the InteractiveServer component. BREAKING CHANGE: Remove NotAvailable variant from MethodStatus enum --- workshop-baker-params/Cargo.lock | 107 +++ workshop-baker-params/Cargo.toml | 8 + workshop-baker-params/src/bake.rs | 24 + workshop-baker-params/src/finalize.rs | 23 + workshop-baker-params/src/lib.rs | 168 ++++ workshop-baker-params/src/notification.rs | 65 ++ workshop-baker-params/src/param.rs | 32 + workshop-baker-params/src/prebake.rs | 54 ++ workshop-baker-params/src/prebake/cache.rs | 37 + workshop-baker-params/src/prebake/engine.rs | 21 + workshop-baker-params/src/prebake/network.rs | 23 + workshop-baker-params/src/traits.rs | 28 + workshop-baker-params/src/writer.rs | 26 + workshop-baker/Cargo.lock | 9 + workshop-baker/Cargo.toml | 1 + workshop-baker/src/finalize/config.rs | 42 + .../src/finalize/plugin/internal.rs | 11 +- .../src/finalize/plugin/metadata.rs | 2 +- workshop-baker/src/main.rs | 8 +- workshop-baker/src/notify.rs | 777 +++++++++--------- workshop-baker/src/notify/interactive.rs | 116 --- workshop-baker/src/notify/types.rs | 23 +- 22 files changed, 1054 insertions(+), 551 deletions(-) create mode 100644 workshop-baker-params/Cargo.lock create mode 100644 workshop-baker-params/Cargo.toml create mode 100644 workshop-baker-params/src/bake.rs create mode 100644 workshop-baker-params/src/finalize.rs create mode 100644 workshop-baker-params/src/lib.rs create mode 100644 workshop-baker-params/src/notification.rs create mode 100644 workshop-baker-params/src/param.rs create mode 100644 workshop-baker-params/src/prebake.rs create mode 100644 workshop-baker-params/src/prebake/cache.rs create mode 100644 workshop-baker-params/src/prebake/engine.rs create mode 100644 workshop-baker-params/src/prebake/network.rs create mode 100644 workshop-baker-params/src/traits.rs create mode 100644 workshop-baker-params/src/writer.rs delete mode 100644 workshop-baker/src/notify/interactive.rs diff --git a/workshop-baker-params/Cargo.lock b/workshop-baker-params/Cargo.lock new file mode 100644 index 0000000..1bff822 --- /dev/null +++ b/workshop-baker-params/Cargo.lock @@ -0,0 +1,107 @@ +# This file is automatically @generated by Cargo. +# It is not intended for manual editing. +version = 4 + +[[package]] +name = "itoa" +version = "1.0.18" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682" + +[[package]] +name = "memchr" +version = "2.8.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "f8ca58f447f06ed17d5fc4043ce1b10dd205e060fb3ce5b979b8ed8e59ff3f79" + +[[package]] +name = "proc-macro2" +version = "1.0.106" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8fd00f0bb2e90d81d1044c2b32617f68fcb9fa3bb7640c23e9c748e53fb30934" +dependencies = [ + "unicode-ident", +] + +[[package]] +name = "quote" +version = "1.0.45" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41f2619966050689382d2b44f664f4bc593e129785a36d6ee376ddf37259b924" +dependencies = [ + "proc-macro2", +] + +[[package]] +name = "serde" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9a8e94ea7f378bd32cbbd37198a4a91436180c5bb472411e48b5ec2e2124ae9e" +dependencies = [ + "serde_core", + "serde_derive", +] + +[[package]] +name = "serde_core" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "41d385c7d4ca58e59fc732af25c3983b67ac852c1a25000afe1175de458b67ad" +dependencies = [ + "serde_derive", +] + +[[package]] +name = "serde_derive" +version = "1.0.228" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "d540f220d3187173da220f885ab66608367b6574e925011a9353e4badda91d79" +dependencies = [ + "proc-macro2", + "quote", + "syn", +] + +[[package]] +name = "serde_json" +version = "1.0.149" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "83fc039473c5595ace860d8c4fafa220ff474b3fc6bfdb4293327f1a37e94d86" +dependencies = [ + "itoa", + "memchr", + "serde", + "serde_core", + "zmij", +] + +[[package]] +name = "syn" +version = "2.0.117" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e665b8803e7b1d2a727f4023456bbbbe74da67099c585258af0ad9c5013b9b99" +dependencies = [ + "proc-macro2", + "quote", + "unicode-ident", +] + +[[package]] +name = "unicode-ident" +version = "1.0.24" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "e6e4313cd5fcd3dad5cafa179702e2b244f760991f45397d14d4ebf38247da75" + +[[package]] +name = "workshop-baker-params" +version = "0.1.0" +dependencies = [ + "serde", + "serde_json", +] + +[[package]] +name = "zmij" +version = "1.0.21" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "b8848ee67ecc8aedbaf3e4122217aff892639231befc6a1b58d29fff4c2cabaa" diff --git a/workshop-baker-params/Cargo.toml b/workshop-baker-params/Cargo.toml new file mode 100644 index 0000000..7819fa9 --- /dev/null +++ b/workshop-baker-params/Cargo.toml @@ -0,0 +1,8 @@ +[package] +name = "workshop-baker-params" +version = "0.1.0" +edition = "2024" + +[dependencies] +serde = { version = "1.0.228", features = ["derive"] } +serde_json = "1.0.149" diff --git a/workshop-baker-params/src/bake.rs b/workshop-baker-params/src/bake.rs new file mode 100644 index 0000000..1eb790d --- /dev/null +++ b/workshop-baker-params/src/bake.rs @@ -0,0 +1,24 @@ +use crate::apply_field; +use serde::{Deserialize, Serialize}; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct BakeParams { + #[serde(default)] pub workspace: ParamVal, +} + +impl MergeParams for BakeParams { + fn defaults() -> Self { Self { workspace: ParamVal::new("/workspace".into()) } } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.workspace, i.workspace, w, ov, p, "workspace"); + } +} + +impl Default for BakeParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct BakeSettings { pub workspace: String } +impl Default for BakeSettings { fn default() -> Self { BakeParams::defaults().into() } } +impl From<&BakeParams> for BakeSettings { fn from(p: &BakeParams) -> Self { Self { workspace: p.workspace.value.clone() } } } +impl From for BakeSettings { fn from(p: BakeParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/finalize.rs b/workshop-baker-params/src/finalize.rs new file mode 100644 index 0000000..15b0578 --- /dev/null +++ b/workshop-baker-params/src/finalize.rs @@ -0,0 +1,23 @@ +use crate::apply_field; +use serde::{Deserialize, Serialize}; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct FinalizeParams { + #[serde(default)] pub artifact_retention_days: ParamVal, +} + +impl MergeParams for FinalizeParams { + fn defaults() -> Self { Self { artifact_retention_days: ParamVal::new(7) } } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.artifact_retention_days, i.artifact_retention_days, w, ov, p, "artifact_retention_days"); + } +} +impl Default for FinalizeParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct FinalizeSettings { pub artifact_retention_days: u32 } +impl Default for FinalizeSettings { fn default() -> Self { FinalizeParams::defaults().into() } } +impl From<&FinalizeParams> for FinalizeSettings { fn from(p: &FinalizeParams) -> Self { Self { artifact_retention_days: p.artifact_retention_days.value } } } +impl From for FinalizeSettings { fn from(p: FinalizeParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/lib.rs b/workshop-baker-params/src/lib.rs new file mode 100644 index 0000000..efff74e --- /dev/null +++ b/workshop-baker-params/src/lib.rs @@ -0,0 +1,168 @@ +pub mod notification; +pub mod prebake; +pub mod bake; +pub mod finalize; +pub mod writer; +pub mod param; +pub mod traits; + +pub use notification::{NotificationParams, NotificationSettings}; +pub use prebake::{PrebakeParams, PrebakeSettings}; +pub use bake::{BakeParams, BakeSettings}; +pub use finalize::{FinalizeParams, FinalizeSettings}; +pub use writer::Writer; +pub use param::ParamVal; +pub use traits::MergeParams; + +use serde::{Deserialize, Serialize}; +use std::collections::HashMap; + +/// Root parameter tree. Each component gets its own subtree. +/// Used for both IPC (Serialized to JSON) and in-process merging. +#[derive(Debug, Clone, Serialize, Deserialize, Default)] +pub struct Params { + #[serde(default, skip_serializing_if = "Option::is_none")] + pub notification: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub prebake: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub bake: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub finalize: Option, +} + +/// Apply and merge helpers. +impl Params { + /// Merge a layer from a specific writer into the current params. + /// Returns entries that were overridden (for debug). + pub fn apply(&mut self, incoming: Params, writer: Writer) -> Overridden { + let mut ov = Overridden::new(); + if let Some(l) = incoming.notification { + if let Some(ref mut b) = self.notification { + b.apply(l, writer, &mut ov, "notification."); + } else { + self.notification = Some(l); + } + } + if let Some(l) = incoming.prebake { + if let Some(ref mut b) = self.prebake { + b.apply(l, writer, &mut ov, "prebake."); + } else { + self.prebake = Some(l); + } + } + if let Some(l) = incoming.bake { + if let Some(ref mut b) = self.bake { + b.apply(l, writer, &mut ov, "bake."); + } else { + self.bake = Some(l); + } + } + if let Some(l) = incoming.finalize { + if let Some(ref mut b) = self.finalize { + b.apply(l, writer, &mut ov, "finalize."); + } else { + self.finalize = Some(l); + } + } + ov + } + + /// Convert to baker-consumable settings. All params guaranteed to have values. + pub fn to_settings(&self) -> Result { + Ok(AllSettings { + notification: self.notification.as_ref().map(NotificationSettings::from).unwrap_or_default(), + prebake: self.prebake.as_ref().map(PrebakeSettings::from).unwrap_or_default(), + bake: self.bake.as_ref().map(BakeSettings::from).unwrap_or_default(), + finalize: self.finalize.as_ref().map(FinalizeSettings::from).unwrap_or_default(), + }) + } +} + +/// Plain-value settings for baker consumption. No metadata, no Option. +#[derive(Debug, Clone)] +pub struct AllSettings { + pub notification: NotificationSettings, + pub prebake: PrebakeSettings, + pub bake: BakeSettings, + pub finalize: FinalizeSettings, +} + +/// Debug — values that were overridden during merge, keyed by dot-path. +#[derive(Debug, Clone, Default)] +pub struct Overridden { + pub entries: std::collections::HashMap, +} + +impl Overridden { + pub fn new() -> Self { Self { entries: HashMap::new() } } + pub fn is_empty(&self) -> bool { self.entries.is_empty() } +} + + + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn test_empty_params_serializes_empty() { + let p = Params::default(); + let json = serde_json::to_string(&p).unwrap(); + assert_eq!(json, "{}"); + } + + #[test] + fn test_notification_params_serialize() { + let mut p = Params::default(); + p.notification = Some(NotificationParams::defaults()); + let json = serde_json::to_string_pretty(&p).unwrap(); + assert!(json.contains("max_lives")); + assert!(json.contains("3")); + } + + #[test] + fn test_apply_courier_overrides_default() { + let mut base = Params::default(); + base.notification = Some(NotificationParams::defaults()); + let incoming = Params { + notification: Some(NotificationParams { + max_lives: ParamVal::with(5, Writer::Courier), + ..Default::default() + }), + ..Default::default() + }; + let ov = base.apply(incoming, Writer::Courier); + let s = base.to_settings().unwrap(); + assert_eq!(s.notification.max_lives, 5); + assert!(ov.is_empty()); // default had no writer, not an "override" + } + + #[test] + fn test_apply_upstream_wins() { + let mut base = Params::default(); + base.notification = Some(NotificationParams { + max_lives: ParamVal::with(3, Writer::Base), + ..Default::default() + }); + let incoming = Params { + notification: Some(NotificationParams { + max_lives: ParamVal::with(5, Writer::Courier), + ..Default::default() + }), + ..Default::default() + }; + let ov = base.apply(incoming, Writer::Courier); + let s = base.to_settings().unwrap(); + // base(3) > courier(2) — rejected, base stays + assert_eq!(s.notification.max_lives, 3); + assert!(ov.is_empty()); // nothing overridden + } + + #[test] + fn test_deny_unknown_fields() { + let bad = r#"{"max_lives":{"value":5},"ghost_param":{"value":1}}"#; + let result: Result = serde_json::from_str(bad); + assert!(result.is_err()); // unknown field should be denied + } +} diff --git a/workshop-baker-params/src/notification.rs b/workshop-baker-params/src/notification.rs new file mode 100644 index 0000000..e4a6e8f --- /dev/null +++ b/workshop-baker-params/src/notification.rs @@ -0,0 +1,65 @@ +use crate::apply_field; +use serde::{Deserialize, Serialize}; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct NotificationParams { + #[serde(default)] pub max_lives: ParamVal, + #[serde(default)] pub timeout_ms: ParamVal, + #[serde(default)] pub retry_backoff: ParamVal, + #[serde(default)] pub retry_delay_seconds: ParamVal, + #[serde(default)] pub max_retry_delay_seconds: ParamVal, + #[serde(default)] pub life_impact_factor: ParamVal, + #[serde(default)] pub batch_period: ParamVal, + #[serde(default)] pub batch_watermark: ParamVal, + #[serde(default)] pub template_ref_depth: ParamVal, +} + +impl MergeParams for NotificationParams { + fn defaults() -> Self { + Self { + max_lives: ParamVal::new(3), timeout_ms: ParamVal::new(30000), + retry_backoff: ParamVal::new(2.0), retry_delay_seconds: ParamVal::new(5.0), + max_retry_delay_seconds: ParamVal::new(300.0), + life_impact_factor: ParamVal::new(0.0), + batch_period: ParamVal::new(30), batch_watermark: ParamVal::new(5), + template_ref_depth: ParamVal::new(10), + } + } + fn apply(&mut self, incoming: Self, writer: Writer, ov: &mut Overridden, prefix: &str) { + apply_field!(self.max_lives, incoming.max_lives, writer, ov, prefix, "max_lives"); + apply_field!(self.timeout_ms, incoming.timeout_ms, writer, ov, prefix, "timeout_ms"); + apply_field!(self.retry_backoff, incoming.retry_backoff, writer, ov, prefix, "retry_backoff"); + apply_field!(self.retry_delay_seconds, incoming.retry_delay_seconds, writer, ov, prefix, "retry_delay_seconds"); + apply_field!(self.max_retry_delay_seconds, incoming.max_retry_delay_seconds, writer, ov, prefix, "max_retry_delay_seconds"); + apply_field!(self.life_impact_factor, incoming.life_impact_factor, writer, ov, prefix, "life_impact_factor"); + apply_field!(self.batch_period, incoming.batch_period, writer, ov, prefix, "batch_period"); + apply_field!(self.batch_watermark, incoming.batch_watermark, writer, ov, prefix, "batch_watermark"); + apply_field!(self.template_ref_depth, incoming.template_ref_depth, writer, ov, prefix, "template_ref_depth"); + } +} + +impl Default for NotificationParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct NotificationSettings { + pub max_lives: u32, pub timeout_ms: u64, pub retry_backoff: f64, + pub retry_delay_seconds: f64, pub max_retry_delay_seconds: f64, + pub life_impact_factor: f64, + pub batch_period: u32, pub batch_watermark: u32, pub template_ref_depth: u32, +} + +impl Default for NotificationSettings { fn default() -> Self { NotificationParams::defaults().into() } } + +impl From<&NotificationParams> for NotificationSettings { + fn from(p: &NotificationParams) -> Self { Self { + max_lives: p.max_lives.value, timeout_ms: p.timeout_ms.value, + retry_backoff: p.retry_backoff.value, retry_delay_seconds: p.retry_delay_seconds.value, + max_retry_delay_seconds: p.max_retry_delay_seconds.value, + life_impact_factor: p.life_impact_factor.value, + batch_period: p.batch_period.value, batch_watermark: p.batch_watermark.value, + template_ref_depth: p.template_ref_depth.value, + }} +} +impl From for NotificationSettings { fn from(p: NotificationParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/param.rs b/workshop-baker-params/src/param.rs new file mode 100644 index 0000000..7b0647a --- /dev/null +++ b/workshop-baker-params/src/param.rs @@ -0,0 +1,32 @@ +use serde::{Deserialize, Serialize}; +use crate::Writer; + +/// A strongly-typed parameter value with writer tracking. +/// +/// Serializes as `{ "value": 3, "writer": "base" }`. +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct ParamVal { + pub value: T, + #[serde(skip_serializing_if = "Option::is_none")] + pub writer: Option, +} + +impl Default for ParamVal { + fn default() -> Self { + ParamVal { value: T::default(), writer: None } + } +} + +impl ParamVal { + pub fn new(value: T) -> Self { + ParamVal { value, writer: None } + } + pub fn with(value: T, writer: Writer) -> Self { + ParamVal { value, writer: Some(writer) } + } + /// Whether an incoming value from `incoming_writer` should override this one. + pub fn should_override(&self, incoming: Writer) -> bool { + let current = self.writer.map_or(0, |w| w.priority()); + incoming.priority() > current + } +} diff --git a/workshop-baker-params/src/prebake.rs b/workshop-baker-params/src/prebake.rs new file mode 100644 index 0000000..53ed9da --- /dev/null +++ b/workshop-baker-params/src/prebake.rs @@ -0,0 +1,54 @@ +pub mod network; +pub mod cache; +pub mod engine; + +use crate::apply_field; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +#[serde(deny_unknown_fields)] +pub struct PrebakeParams { + #[serde(default)] pub bootstrap_user: ParamVal, + #[serde(default)] pub workspace: ParamVal, + #[serde(default)] pub network: network::NetworkParams, + #[serde(default)] pub cache: cache::CacheParams, + #[serde(default)] pub engine: engine::EngineParams, + #[serde(default)] pub repology_endpoint: ParamVal, +} + +impl MergeParams for PrebakeParams { + fn defaults() -> Self { + Self { + bootstrap_user: ParamVal::new("vulcan".into()), workspace: ParamVal::new("/home/vulcan/workspace".into()), + network: network::NetworkParams::defaults(), cache: cache::CacheParams::defaults(), + engine: engine::EngineParams::defaults(), + repology_endpoint: ParamVal::new("default".into()), + } + } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.bootstrap_user, i.bootstrap_user, w, ov, p, "bootstrap_user"); + apply_field!(self.workspace, i.workspace, w, ov, p, "workspace"); + apply_field!(self.repology_endpoint, i.repology_endpoint, w, ov, p, "repology_endpoint"); + self.network.apply(i.network, w, ov, &format!("{}network.", p)); + self.cache.apply(i.cache, w, ov, &format!("{}cache.", p)); + self.engine.apply(i.engine, w, ov, &format!("{}engine.", p)); + } +} +impl Default for PrebakeParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct PrebakeSettings { + pub bootstrap_user: String, pub workspace: String, + pub network: network::NetworkSettings, pub cache: cache::CacheSettings, + pub engine: engine::EngineSettings, pub repology_endpoint: String, +} +impl Default for PrebakeSettings { fn default() -> Self { PrebakeParams::defaults().into() } } +impl From<&PrebakeParams> for PrebakeSettings { + fn from(p: &PrebakeParams) -> Self { Self { + bootstrap_user: p.bootstrap_user.value.clone(), workspace: p.workspace.value.clone(), + network: (&p.network).into(), cache: (&p.cache).into(), engine: (&p.engine).into(), + repology_endpoint: p.repology_endpoint.value.clone(), + }} +} +impl From for PrebakeSettings { fn from(p: PrebakeParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/prebake/cache.rs b/workshop-baker-params/src/prebake/cache.rs new file mode 100644 index 0000000..f72058f --- /dev/null +++ b/workshop-baker-params/src/prebake/cache.rs @@ -0,0 +1,37 @@ +use crate::apply_field; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct CacheParams { + #[serde(default)] pub strategy: ParamVal, + #[serde(default)] pub mode: ParamVal, + #[serde(default)] pub ttl_days: ParamVal, + #[serde(default)] pub max_size_gb: ParamVal, + #[serde(default)] pub cleanup_policy: ParamVal, +} +impl MergeParams for CacheParams { + fn defaults() -> Self { + Self { + strategy: ParamVal::new("always".into()), mode: ParamVal::new("zstd".into()), + ttl_days: ParamVal::new(30), max_size_gb: ParamVal::new(20), + cleanup_policy: ParamVal::new("lru".into()), + } + } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.strategy, i.strategy, w, ov, p, "strategy"); + apply_field!(self.mode, i.mode, w, ov, p, "mode"); + apply_field!(self.ttl_days, i.ttl_days, w, ov, p, "ttl_days"); + apply_field!(self.max_size_gb, i.max_size_gb, w, ov, p, "max_size_gb"); + apply_field!(self.cleanup_policy, i.cleanup_policy, w, ov, p, "cleanup_policy"); + } +} +impl Default for CacheParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct CacheSettings { pub strategy: String, pub mode: String, pub ttl_days: u32, pub max_size_gb: u32, pub cleanup_policy: String } +impl Default for CacheSettings { fn default() -> Self { CacheParams::defaults().into() } } +impl From<&CacheParams> for CacheSettings { fn from(p: &CacheParams) -> Self { + Self { strategy: p.strategy.value.clone(), mode: p.mode.value.clone(), ttl_days: p.ttl_days.value, max_size_gb: p.max_size_gb.value, cleanup_policy: p.cleanup_policy.value.clone() } +} } +impl From for CacheSettings { fn from(p: CacheParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/prebake/engine.rs b/workshop-baker-params/src/prebake/engine.rs new file mode 100644 index 0000000..c1a487e --- /dev/null +++ b/workshop-baker-params/src/prebake/engine.rs @@ -0,0 +1,21 @@ +use crate::apply_field; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct EngineParams { + #[serde(default)] pub timeout_ms: ParamVal, +} +impl MergeParams for EngineParams { + fn defaults() -> Self { Self { timeout_ms: ParamVal::new(300000) } } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.timeout_ms, i.timeout_ms, w, ov, p, "timeout_ms"); + } +} +impl Default for EngineParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct EngineSettings { pub timeout_ms: u64 } +impl Default for EngineSettings { fn default() -> Self { EngineParams::defaults().into() } } +impl From<&EngineParams> for EngineSettings { fn from(p: &EngineParams) -> Self { Self { timeout_ms: p.timeout_ms.value } } } +impl From for EngineSettings { fn from(p: EngineParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/prebake/network.rs b/workshop-baker-params/src/prebake/network.rs new file mode 100644 index 0000000..42f7146 --- /dev/null +++ b/workshop-baker-params/src/prebake/network.rs @@ -0,0 +1,23 @@ +use crate::apply_field; +use crate::{ParamVal, Writer, Overridden, traits::MergeParams}; +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Serialize, Deserialize)] +pub struct NetworkParams { + #[serde(default)] pub enabled: ParamVal, + #[serde(default)] pub outbound: ParamVal, +} +impl MergeParams for NetworkParams { + fn defaults() -> Self { Self { enabled: ParamVal::new(true), outbound: ParamVal::new(true) } } + fn apply(&mut self, i: Self, w: Writer, ov: &mut Overridden, p: &str) { + apply_field!(self.enabled, i.enabled, w, ov, p, "enabled"); + apply_field!(self.outbound, i.outbound, w, ov, p, "outbound"); + } +} +impl Default for NetworkParams { fn default() -> Self { Self::defaults() } } + +#[derive(Debug, Clone)] +pub struct NetworkSettings { pub enabled: bool, pub outbound: bool } +impl Default for NetworkSettings { fn default() -> Self { NetworkParams::defaults().into() } } +impl From<&NetworkParams> for NetworkSettings { fn from(p: &NetworkParams) -> Self { Self { enabled: p.enabled.value, outbound: p.outbound.value } } } +impl From for NetworkSettings { fn from(p: NetworkParams) -> Self { (&p).into() } } diff --git a/workshop-baker-params/src/traits.rs b/workshop-baker-params/src/traits.rs new file mode 100644 index 0000000..ddb893b --- /dev/null +++ b/workshop-baker-params/src/traits.rs @@ -0,0 +1,28 @@ +use crate::{Writer, Overridden}; + +/// A parameter tree node that can merge itself with an incoming layer. +pub trait MergeParams: Sized { + fn apply(&mut self, incoming: Self, writer: Writer, ov: &mut Overridden, prefix: &str); + fn defaults() -> Self; +} + +/// Apply one field's merge logic. +/// ```ignore +/// apply_field!(self.foo, incoming.foo, writer, ov, prefix, "foo"); +/// ``` +#[macro_export] +macro_rules! apply_field { + ($self:ident.$field:ident, $incoming:ident.$field2:ident, $writer:expr, $ov:ident, $prefix:expr, $name:expr) => { + let old_w = $self.$field.writer; + if $self.$field.should_override($writer) { + let old_v = ::std::mem::replace(&mut $self.$field.value, $incoming.$field2.value); + $self.$field.writer = ::std::option::Option::Some($writer); + if let ::std::option::Option::Some(w) = old_w { + $ov.entries.insert( + format!("{}{}", $prefix, $name), + ::serde_json::json!({"value": old_v, "writer": w.as_str()}), + ); + } + } + }; +} diff --git a/workshop-baker-params/src/writer.rs b/workshop-baker-params/src/writer.rs new file mode 100644 index 0000000..c3169d3 --- /dev/null +++ b/workshop-baker-params/src/writer.rs @@ -0,0 +1,26 @@ +use serde::{Deserialize, Serialize}; + +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "lowercase")] +pub enum Writer { + Base, + Courier, + Bakerd, +} + +impl Writer { + /// Priority: Base=3 > Courier=2 > Bakerd=1 + pub fn priority(self) -> u8 { + match self { + Writer::Base => 3, + Writer::Courier => 2, + Writer::Bakerd => 1, + } + } + pub fn as_str(self) -> &'static str { + match self { + Writer::Base => "base", Writer::Courier => "courier", Writer::Bakerd => "bakerd", + } + } +} +impl Default for Writer { fn default() -> Self { Writer::Base } } diff --git a/workshop-baker/Cargo.lock b/workshop-baker/Cargo.lock index 1fe400f..d0e5c1b 100644 --- a/workshop-baker/Cargo.lock +++ b/workshop-baker/Cargo.lock @@ -3847,9 +3847,18 @@ dependencies = [ "topological-sort", "uuid", "which", + "workshop-baker-params", "workshop-engine", ] +[[package]] +name = "workshop-baker-params" +version = "0.1.0" +dependencies = [ + "serde", + "serde_json", +] + [[package]] name = "workshop-engine" version = "0.1.0" diff --git a/workshop-baker/Cargo.toml b/workshop-baker/Cargo.toml index c3413ed..f336791 100644 --- a/workshop-baker/Cargo.toml +++ b/workshop-baker/Cargo.toml @@ -50,6 +50,7 @@ privdrop = "0.5.6" which = "8.0.2" globset = "0.4.18" axum = "0.7" +workshop-baker-params = { version = "0.1.0", path = "../workshop-baker-params" } [dev-dependencies] criterion = "0.5" diff --git a/workshop-baker/src/finalize/config.rs b/workshop-baker/src/finalize/config.rs index ec706c6..ec0359b 100644 --- a/workshop-baker/src/finalize/config.rs +++ b/workshop-baker/src/finalize/config.rs @@ -115,6 +115,17 @@ pub struct NotificationMethod { pub policy: Vec, } +impl Default for NotificationMethod { + fn default() -> Self { + Self { + template: NotificationTemplate::default(), + trigger: Vec::new(), + config: Value::default(), + policy: Vec::new(), + } + } +} + #[derive(Debug, Deserialize, Serialize, Clone, Default)] pub struct NotificationTemplate { /// Payload schema: "default" (plugin-defined) or "custom" (requires on_* templates) @@ -134,6 +145,23 @@ pub struct NotificationTemplate { pub on_finish_of: Option, } +impl std::ops::AddAssign for NotificationTemplate { + fn add_assign(&mut self, other: Self) { + if self.schema.is_none() { + self.schema = other.schema; + } + if self.on_success.is_none() { + self.on_success = other.on_success; + } + if self.on_failure.is_none() { + self.on_failure = other.on_failure; + } + if self.on_finish_of.is_none() { + self.on_finish_of = other.on_finish_of; + } + } +} + /// Notification configuration container #[derive(Debug, Serialize, Clone, Default)] pub struct NotificationConfig { @@ -203,6 +231,20 @@ pub struct NotificationTemplateDef { pub on_finish_of: Option, } +impl std::ops::AddAssign for NotificationTemplateDef { + fn add_assign(&mut self, other: Self) { + if self.on_success.is_none() { + self.on_success = other.on_success; + } + if self.on_failure.is_none() { + self.on_failure = other.on_failure; + } + if self.on_finish_of.is_none() { + self.on_finish_of = other.on_finish_of; + } + } +} + /// Notification trigger type. #[derive(Debug, Clone, PartialEq, Eq)] pub enum Trigger { diff --git a/workshop-baker/src/finalize/plugin/internal.rs b/workshop-baker/src/finalize/plugin/internal.rs index 8e54751..c537499 100644 --- a/workshop-baker/src/finalize/plugin/internal.rs +++ b/workshop-baker/src/finalize/plugin/internal.rs @@ -27,9 +27,12 @@ pub fn get_internal_plugin() -> Vec { name: "webhook", entry: Box::new(webhook::Webhook), }, - InternalPlugin { - name: "dummy", - entry: Box::new(dummy::Dummy), - }, ] } + +pub fn get_dummy_plugin() -> Vec { + vec![InternalPlugin { + name: "dummy", + entry: Box::new(dummy::Dummy), + }] +} diff --git a/workshop-baker/src/finalize/plugin/metadata.rs b/workshop-baker/src/finalize/plugin/metadata.rs index 68ac70b..24fea29 100644 --- a/workshop-baker/src/finalize/plugin/metadata.rs +++ b/workshop-baker/src/finalize/plugin/metadata.rs @@ -19,7 +19,7 @@ pub struct PluginMetadata { #[serde(default)] pub renderer: Option, - /// Additional YAML fields not captured by struct. + /// Additional YAML fields. #[serde(flatten)] pub value: serde_yaml::Value, diff --git a/workshop-baker/src/main.rs b/workshop-baker/src/main.rs index 2408a59..d910b9e 100644 --- a/workshop-baker/src/main.rs +++ b/workshop-baker/src/main.rs @@ -16,15 +16,15 @@ fn or_workshop(path: &Option, name: &str) -> PathBuf { async fn worker_main(cli: Cli) -> Result<(), CliError> { let prebake_path = cli.prebake.clone().unwrap_or_else(|| { - eprintln!("Error: --prebake is required in worker standalone mode"); + log::error!("--prebake is required in worker standalone mode"); std::process::exit(1); }); let bake_path = cli.bake.clone().unwrap_or_else(|| { - eprintln!("Error: --bake is required in worker standalone mode"); + log::error!("--bake is required in worker standalone mode"); std::process::exit(1); }); let finalize_path = cli.finalize.clone().unwrap_or_else(|| { - eprintln!("Error: --finalize is required in worker standalone mode"); + log::error!("--finalize is required in worker standalone mode"); std::process::exit(1); }); @@ -59,7 +59,7 @@ async fn worker_main(cli: Cli) -> Result<(), CliError> { let debug = cli.debug.clone(); let notify_registry = registry.clone(); let notify_handle = tokio::spawn(async move { - notify(¬ify_path, &mut notify_ctx.clone(), notify_etx, to_rx, retry_tx, debug, ¬ify_registry).await + notify(¬ify_path, &mut notify_ctx.clone(), notify_etx, to_rx, retry_tx, debug.map(|d| d.contains("notify")).unwrap_or(false), ¬ify_registry).await }); let nevent_tx = to_tx.clone(); diff --git a/workshop-baker/src/notify.rs b/workshop-baker/src/notify.rs index 34cd43c..65c7a4e 100644 --- a/workshop-baker/src/notify.rs +++ b/workshop-baker/src/notify.rs @@ -2,21 +2,12 @@ use crate::finalize::FinalizeConfig; use crate::finalize::NotificationConfig; use crate::finalize::NotificationTemplate; use crate::finalize::config::TriggerItemRaw; -use crate::types::resource::ResourceRegistry; use crate::finalize::parse; use crate::finalize::plugin; use crate::notify::types::Matchable; use crate::notify::types::Method; use crate::notify::types::MethodStatus; -use crate::notify::types::NOTIFICATION_DEFAULT_FALLIBLE; -use crate::notify::types::NOTIFICATION_DEFAULT_MAX_LIVES; -use crate::notify::types::NOTIFICATION_DEFAULT_RETRY_BACKOFF; -use crate::notify::types::NOTIFICATION_DEFAULT_TIMEOUT_MS; -use crate::notify::types::NOTIFICATION_LIFE_IMPACT_FACTOR; -use crate::notify::types::NOTIFICATION_NOTREADY_DELAY; -use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_FACTOR; -use crate::notify::types::NOTIFICATION_TEMPORARY_DELAY_TIME; -use crate::notify::types::NOTIFICATION_TEMPORARY_MAX_TIME; +use crate::notify::types::DEFAULT_NOTIFY_PARAMS; use crate::notify::types::NotificationEvent; use crate::notify::types::NotificationEventReceiver; use crate::notify::types::NotificationEventSender; @@ -25,18 +16,18 @@ use crate::notify::types::NotificationRenderer; use crate::notify::types::NotificationRule; use crate::notify::types::NotificationRulesMap; use crate::notify::types::ParsedTrigger; -use workshop_engine::{EventSender, ExecutionContext}; +use crate::types::resource::ResourceRegistry; use chrono::Utc; use globset::{Glob, GlobSetBuilder}; use std::collections::HashMap; use std::collections::HashSet; use std::path::Path; +use workshop_engine::{EventSender, ExecutionContext}; mod error; -pub mod types; -pub mod template; -pub mod interactive; pub mod queue; +pub mod template; +pub mod types; use error::NotifyError; /// Maximum depth for template cross-references (prevent infinite recursion) @@ -69,29 +60,22 @@ pub async fn notify( _event_tx: EventSender, mut rx: NotificationEventReceiver, retry_tx: NotificationEventSender, - debug: Option, + debug: bool, registry: &ResourceRegistry, ) -> anyhow::Result<()> { let mut finalize = parse(finalize_path)?; validate(&finalize)?; let mut plugin_list = extract_plugin(&finalize); - let is_debug = debug.as_deref().unwrap_or("").contains("notify"); - if plugin_list.is_empty() { - if is_debug { - log::info!("[NOTIFY_DEBUG] Injecting dummy plugin for debug observation"); + if debug { + log::debug!("[NOTIFY_DEBUG] Injecting dummy plugin."); plugin_list.insert("dummy".to_string()); // Insert a dummy method so the notification system has at least one entry let mut methods = std::collections::HashMap::new(); methods.insert( "dummy".to_string(), - crate::finalize::NotificationMethod { - template: crate::finalize::NotificationTemplate::default(), - trigger: vec![], - config: serde_yaml::Value::Mapping(serde_yaml::Mapping::new()), - policy: vec![], - }, + crate::finalize::NotificationMethod::default(), ); finalize.notification.groups.insert(0, methods); } else { @@ -100,7 +84,12 @@ pub async fn notify( } } - let notify_plugins = plugin::register(finalize.plugin, Some(&plugin_list), &mut finalize.notification, registry)?; + let notify_plugins = plugin::register( + finalize.plugin, + Some(&plugin_list), + &mut finalize.notification, + registry, + )?; // Resolve external template references from registry (pre-fetched) for (name, tmpl_def) in finalize.notification.templates.iter_mut() { @@ -109,10 +98,13 @@ pub async fn notify( let content = std::fs::read_to_string(path.join("template.yml")) .or_else(|_| std::fs::read_to_string(path.join("template.yaml"))); if let Ok(content) = content { - if let Ok(fetched) = serde_yaml::from_str::(&content) { + if let Ok(fetched) = + serde_yaml::from_str::(&content) + { tmpl_def.on_success = fetched.on_success.or(tmpl_def.on_success.take()); tmpl_def.on_failure = fetched.on_failure.or(tmpl_def.on_failure.take()); - tmpl_def.on_finish_of = fetched.on_finish_of.or(tmpl_def.on_finish_of.take()); + tmpl_def.on_finish_of = + fetched.on_finish_of.or(tmpl_def.on_finish_of.take()); tmpl_def.use_ = None; // Resolved } } @@ -125,24 +117,11 @@ pub async fn notify( let notification_rules = parse_rules(&finalize.notification); - loop { // Wait for first event - let first = if is_debug { - // Debug mode: don't block forever — timeout prevents deadlock - match tokio::time::timeout(std::time::Duration::from_secs(5), rx.recv()).await { - Ok(Some(event)) => event, - Ok(None) => break, - Err(_) => { - log::info!("[NOTIFY_DEBUG] No events within timeout, shutting down"); - break; - } - } - } else { - match rx.recv().await { - Some(event) => event, - None => break, - } + let first = match rx.recv().await { + Some(event) => event, + None => break, }; let mut events = vec![first]; @@ -151,48 +130,52 @@ pub async fn notify( events.push(event); } - if debug.as_deref().unwrap_or("").contains("notify") { + if debug { for ev in &events { - log::info!("[NOTIFY_DEBUG] stage={:?} substage={} outcome={:?} priority={}", - ev.stage, ev.substage, ev.outcome, ev.priority); + log::info!( + "[NOTIFY_DEBUG] stage={:?} substage={} outcome={:?} priority={}", + ev.stage, + ev.substage, + ev.outcome, + ev.priority + ); } } - let results = notification_handler(¬ify_plugins, ¬ification_rules, &events, &finalize.notification.templates, ctx).await; + let results = notification_handler( + ¬ify_plugins, + ¬ification_rules, + &events, + &finalize.notification.templates, + ctx, + ) + .await; - let mut na: Vec = vec![]; - let mut te: Vec = vec![]; - let mut pe_ex: Vec<(String, bool)> = vec![]; + let mut recoverable: Vec = vec![]; + let mut unrecoverable: Vec<(String, bool)> = vec![]; for (method, status) in &results { match status { MethodStatus::Success => {} - MethodStatus::NotAvailable(_) => na.push(method.clone()), - MethodStatus::TemporaryError(_) => te.push(method.clone()), + MethodStatus::TemporaryError(_) => recoverable.push(method.clone()), MethodStatus::PermanentError(_) | MethodStatus::Exhausted(_) => { let fallible = get_fallible(¬ification_rules, method); - pe_ex.push((method.clone(), fallible)); + unrecoverable.push((method.clone(), fallible)); } } } - if !na.is_empty() || !te.is_empty() { - let mut retry_methods = na.clone(); - retry_methods.extend(te.clone()); + if !recoverable.is_empty() { + let retry_methods = recoverable.clone(); events[0].target.rules = Some(retry_methods); - if !te.is_empty() { - events[0].target.count += 1; - events[0].not_before = Utc::now() + notify_delay(&events[0]); - } else { - events[0].effective_priority += 1; - events[0].not_before = Utc::now() + NOTIFICATION_NOTREADY_DELAY; - } + events[0].target.count += 1; + events[0].not_before = Utc::now() + notify_delay(&events[0]); retry_tx.send(events[0].clone()).await?; continue; } - let has_hard = pe_ex.iter().any(|(_, fallible)| !fallible); - if has_hard { + let failed = unrecoverable.iter().any(|(_, fallible)| !fallible); + if failed { events[0].target.priority += 1; events[0].target.rules = None; events[0].attempted = events[0].target.count; @@ -207,85 +190,95 @@ pub async fn notify( // TODO: PIPELINE_NOTIFICATION_EXHAUSTED — set when any method returns Exhausted // TODO: PIPELINE_NOTIFICATION_FAILED — set when any method returns PermanentError -// Markers sent via event_tx as ExecutionEvent variants for WebUI/CLI display. +// Markers sent via event_tx as ExecutionEvent variants for future WebUI/CLI display. - -/// Resolve a schema reference into a concrete NotificationTemplateDef. +/// Resolve a schema reference into a NotificationTemplateDef by walking a fallback chain. /// -/// Resolution rules: -/// - `None` (unspecified): c/d → ws-base default → fallback; standalone → fallback -/// - `Some("custom")`: return empty default (caller uses rule.template.on_* directly) -/// - `Some("default")`: ws-base default (fallback in standalone) -/// - `Some("")`: look up in templates; if not found and name == method_name, -/// fallback to `{name}.fallback`; otherwise error -fn resolve_template_content( +/// The chain is defined by resolve_template_next_schema(): each entry point prescribes +/// a sequence of schema names to try. Each step calls resolve_template_by_schema() to do +/// the actual lookup; the first hit is returned as Ok(...). If the chain is exhausted +/// without a hit, a descriptive SchemaError is returned. +/// +/// Behaviour difference between standalone and daemon modes is implicit in +/// server_templates (None in standalone, Some in daemon). No explicit mode check needed. +fn resolve_template( schema: &Option, method_name: &str, templates: &std::collections::HashMap, - standalone: bool, - ws_base_templates: Option<&std::collections::HashMap>, + _standalone: bool, + server_templates: Option< + &std::collections::HashMap, + >, ) -> Result { + // 1. Try the initial schema directly + if let Some(schema_name) = schema.as_deref() { + if let Some(template) = + resolve_template_by_schema(schema_name, method_name, templates, server_templates) + { + return Ok(template); + } + } + + // 2. Walk the fallback chain + let mut schema_next = resolve_template_next_schema(schema); + while let Some(ref schema_name) = schema_next { + if let Some(template) = + resolve_template_by_schema(schema_name, method_name, templates, server_templates) + { + return Ok(template); + } + log::warn!("Failed to resolve template for schema {}.", schema_name); + schema_next = resolve_template_next_schema(&schema_next); + } + Err(NotifyError::SchemaError(format!("All template resolution attempts beginning from {} failed.", schema.as_deref().unwrap_or("[system default]")))) +} + +/// Return the next schema name in the fallback chain, or None to terminate. +/// +/// This is the transition function of a state machine driven by resolve_template(). +/// Each state (schema name) maps to a next attempt: +/// +/// | current schema | next schema | why | +/// |---|---:|------| +/// | `"custom"` | None | Schema says "plugin does its own rendering". Not really a template resolution. | +/// | `"default"` | `"fallback"` | System default not found → try plugin's own built-in template. | +/// | `"fallback"` | None | Last resort exhausted. | +/// | `""` | `"default"` | User-named template not found → fall back to system default. | +/// | `None` (unset) | `"default"` | No explicit schema → start from system default. | +fn resolve_template_next_schema(schema: &Option) -> Option { match schema.as_deref() { - Some("custom") => { - // Caller reads template content from rule.template.on_* directly - Ok(crate::finalize::NotificationTemplateDef::default()) - } + Some("custom") => None, + Some("default") => Some("fallback".to_string()), + Some("fallback") => None, + Some(_) => Some("default".to_string()), + None => Some("default".to_string()), + } +} - Some("default") => { - if standalone { - templates.get(&format!("{}.fallback", method_name)) - .cloned() - .ok_or_else(|| NotifyError::SchemaError( - format!("{} fallback template not found in standalone mode", method_name) - )) - } else { - ws_base_templates - .and_then(|b| b.get("default")) - .cloned() - .ok_or_else(|| NotifyError::SchemaError( - "ws-base default template not found".into() - )) - } - } - - Some(name) => { - match templates.get(name) { - Some(tmpl) => Ok(tmpl.clone()), - None => { - if name == method_name { - // Fallback to plugin's own template - templates.get(&format!("{}.fallback", name)) - .cloned() - .ok_or_else(|| NotifyError::SchemaError( - format!("plugin '{}' has no fallback template", name) - )) - } else { - Err(NotifyError::SchemaError( - format!("template '{}' not found", name) - )) - } - } - } - } - - None => { - if standalone { - templates.get(&format!("{}.fallback", method_name)) - .cloned() - .ok_or_else(|| NotifyError::SchemaError( - format!("no template configured and '{}' has no fallback", method_name) - )) - } else { - // c/d: ws-base default → fallback (two-level chain) - ws_base_templates - .and_then(|base| base.get("default")) - .or_else(|| templates.get(&format!("{}.fallback", method_name))) - .cloned() - .ok_or_else(|| NotifyError::SchemaError( - "no template available (ws-base default not found, no plugin fallback)".into() - )) - } - } +/// Perform a single lookup for a given schema name. Returns None if not found. +/// +/// This is the action function of the state machine — it has no knowledge of the +/// fallback chain and leaves progression to resolve_template(). +/// +/// | schema | lookup target | returns | +/// |--------------|-----------------------------------------------|---------| +/// | `"custom"` | — | empty default (plugin self-renders) | +/// | `"fallback"` | templates[`"{method}.fallback"`] | plugin's built-in template | +/// | `"default"` | server_templates[method_name] | bakerd's system default | +/// | other name | templates[name] | user-defined template | +fn resolve_template_by_schema( + schema: &str, + method_name: &str, + templates: &std::collections::HashMap, + server_templates: Option< + &std::collections::HashMap, + >, +) -> Option { + match schema { + "custom" => Some(crate::finalize::NotificationTemplateDef::default()), + "fallback" => templates.get(&format!("{}.fallback", method_name)).cloned(), + "default" => server_templates.and_then(|b| b.get(method_name)).cloned(), + name => templates.get(name).cloned(), } } @@ -334,11 +327,9 @@ fn resolve_chain( // Merge: base provides defaults, referencing template overrides let base = templates[&ref_name].clone(); - let tmpl = templates.get_mut(name).unwrap(); - if tmpl.on_success.is_none() { tmpl.on_success = base.on_success.clone(); } - if tmpl.on_failure.is_none() { tmpl.on_failure = base.on_failure.clone(); } - if tmpl.on_finish_of.is_none() { tmpl.on_finish_of = base.on_finish_of.clone(); } - tmpl.use_ = None; // Chain flattened + let template = templates.get_mut(name).unwrap(); + *template += base; + template.use_ = None; // Chain flattened } } @@ -346,10 +337,6 @@ fn resolve_chain( Ok(()) } -fn is_plugin_unavailable(e: &crate::finalize::error::PluginError) -> bool { - matches!(e, crate::finalize::error::PluginError::PluginUnavailable { .. }) -} - fn is_temporary_error(e: &crate::finalize::error::PluginError) -> bool { !matches!(e, crate::finalize::error::PluginError::GeneralError(_)) } @@ -362,7 +349,11 @@ async fn notification_handler( ctx: &ExecutionContext, ) -> Vec<(String, MethodStatus)> { let mut results: Vec<(String, MethodStatus)> = vec![]; - log::info!("[NOTIFY] handler called with {} events, {} priority groups", events.len(), config.len()); + log::info!( + "[NOTIFY] handler called with {} events, {} priority groups", + events.len(), + config.len() + ); let mut buckets: std::collections::HashMap<(String, usize), Vec<&NotificationEvent>> = std::collections::HashMap::new(); @@ -376,15 +367,27 @@ async fn notification_handler( } } for (rule_id, rule) in rules.iter().enumerate() { - log::info!("[NOTIFY] rule_id={} enabled={} trigger_len={}", rule_id, rule.enabled, rule.trigger.len()); - eprintln!("[CHK-ENABLED] entering enabled check"); + log::info!( + "[NOTIFY] rule_id={} enabled={} trigger_len={}", + rule_id, + rule.enabled, + rule.trigger.len() + ); + log::trace!("[CHK-ENABLED] entering enabled check"); if !rule.enabled { - eprintln!("[CHK-ENABLED] enabled=false, continuing"); + log::trace!("[CHK-ENABLED] enabled=false, continuing"); continue; } - eprintln!("[BEFORE-ACCEPT] method={} rule={} event={}.{}", method, rule_id, event.stage, event.substage); + log::trace!("[BEFORE-ACCEPT] method={} rule={} event={}.{}", + method, rule_id, event.stage, event.substage); if rule.accept(&event.outcome, &event.stage, &event.substage) { - log::info!("[NOTIFY] MATCH: method={} rule={} event={}.{}", method, rule_id, event.stage, event.substage); + log::info!( + "[NOTIFY] MATCH: method={} rule={} event={}.{}", + method, + rule_id, + event.stage, + event.substage + ); buckets .entry((method.clone(), rule_id)) .or_default() @@ -404,7 +407,10 @@ async fn notification_handler( let rule = match find_rule(config, &method, rule_id) { Ok(r) => r, Err(e) => { - results.push((method.clone(), MethodStatus::PermanentError(format!("{}", e)))); + results.push(( + method.clone(), + MethodStatus::PermanentError(format!("{}", e)), + )); continue; } }; @@ -412,7 +418,10 @@ async fn notification_handler( let plugin = match plugins.get(&method) { Some(p) => p, None => { - results.push((method.clone(), MethodStatus::PermanentError(format!("Unknown plugin: {}", method)))); + results.push(( + method.clone(), + MethodStatus::PermanentError(format!("Unknown plugin: {}", method)), + )); continue; } }; @@ -422,7 +431,7 @@ async fn notification_handler( let argument = match renderer { NotificationRenderer::Server => { - let template_def = match resolve_template_content( + let template_def = match resolve_template( &rule.template.schema, &method, templates, @@ -431,19 +440,35 @@ async fn notification_handler( ) { Ok(td) => td, Err(e) => { - results.push((method.clone(), MethodStatus::PermanentError(format!("{}", e)))); + results.push(( + method.clone(), + MethodStatus::PermanentError(format!("{}", e)), + )); continue; } }; let first = &bucket[0]; - let template_str = match crate::notify::template::select_template_str(&template_def, &first.outcome) - .or_else(|| rule.template.on_success.as_deref()) - .or_else(|| rule.template.on_failure.as_deref()) - .or_else(|| rule.template.on_finish_of.as_ref().map(|o| o.content.as_str())) - { + let template_str = match crate::notify::template::select_template_str( + &template_def, + &first.outcome, + ) + .or_else(|| rule.template.on_success.as_deref()) + .or_else(|| rule.template.on_failure.as_deref()) + .or_else(|| { + rule.template + .on_finish_of + .as_ref() + .map(|o| o.content.as_str()) + }) { Some(ts) => ts, None => { - results.push((method.clone(), MethodStatus::PermanentError(format!("no template for outcome {:?}", first.outcome)))); + results.push(( + method.clone(), + MethodStatus::PermanentError(format!( + "no template for outcome {:?}", + first.outcome + )), + )); continue; } }; @@ -471,7 +496,7 @@ async fn notification_handler( NotificationRenderer::Plugin => { let mut argument = rule.config.clone(); - if let Ok(template_def) = resolve_template_content( + if let Ok(template_def) = resolve_template( &rule.template.schema, &method, templates, @@ -480,10 +505,7 @@ async fn notification_handler( ) { if let serde_yaml::Value::Mapping(ref mut map) = argument { if let Ok(tmpl_value) = serde_yaml::to_value(&template_def) { - map.insert( - serde_yaml::Value::String("_template".into()), - tmpl_value, - ); + map.insert(serde_yaml::Value::String("_template".into()), tmpl_value); } } } @@ -511,66 +533,20 @@ async fn notification_handler( } NotificationRenderer::Interactive => { - let template_data = - crate::notify::template::build_batch_template_context(ctx, &bucket); - let mut server = - crate::notify::interactive::InteractiveServer::new(template_data); - let port = match server.start().await { - Ok(p) => p, - Err(e) => { - results.push((method.clone(), MethodStatus::TemporaryError(format!("Failed to start IR server: {}", e)))); - continue; - } - }; - - let mut argument = rule.config.clone(); - if let serde_yaml::Value::Mapping(ref mut map) = argument { - map.insert( - serde_yaml::Value::String("HBW_IR_PORT".into()), - serde_yaml::Value::Number(port.into()), - ); - } - - log::info!( - "IR: server started on port {} for method '{}'", - port, method - ); - - let remain = rule.max_lives as f64 - - NOTIFICATION_LIFE_IMPACT_FACTOR * events[0].attempted as f64 - - events[0].target.count as f64; - - if remain < 1.0 { - results.push((method.clone(), MethodStatus::Exhausted(format!("max_lives depleted")))); - server.shutdown().await; - continue; - } - - let status = match tokio::time::timeout(rule.timeout, plugin.call(argument, ctx, &None)).await { - Ok(Ok(())) => MethodStatus::Success, - Ok(Err(e)) => { - if is_plugin_unavailable(&e) { - MethodStatus::NotAvailable(format!("{}", e)) - } else if is_temporary_error(&e) { - MethodStatus::TemporaryError(format!("{}", e)) - } else { - MethodStatus::PermanentError(format!("{}", e)) - } - } - Err(_) => MethodStatus::TemporaryError("timeout".into()), - }; - server.shutdown().await; - results.push((method.clone(), status)); - continue; + unimplemented!("Interactive renderer is not yet implemented") } + }; let remain = rule.max_lives as f64 - - NOTIFICATION_LIFE_IMPACT_FACTOR * events[0].attempted as f64 + - DEFAULT_NOTIFY_PARAMS.life_impact_factor.value * events[0].attempted as f64 - events[0].target.count as f64; if remain < 1.0 { - results.push((method.clone(), MethodStatus::Exhausted(format!("max_lives depleted")))); + results.push(( + method.clone(), + MethodStatus::Exhausted(format!("max_lives depleted")), + )); continue; } @@ -579,9 +555,7 @@ async fn notification_handler( results.push((method.clone(), MethodStatus::Success)); } Ok(Err(e)) => { - let status = if is_plugin_unavailable(&e) { - MethodStatus::NotAvailable(format!("{}", e)) - } else if is_temporary_error(&e) { + let status = if is_temporary_error(&e) { MethodStatus::TemporaryError(format!("{}", e)) } else { MethodStatus::PermanentError(format!("{}", e)) @@ -589,7 +563,10 @@ async fn notification_handler( results.push((method.clone(), status)); } Err(_) => { - results.push((method.clone(), MethodStatus::TemporaryError("timeout".into()))); + results.push(( + method.clone(), + MethodStatus::TemporaryError("timeout".into()), + )); } } } @@ -597,45 +574,52 @@ async fn notification_handler( results } +/// Parse a single `on_finish_of` pattern string into a ParsedTrigger, or None. +/// Supports `re:` prefix for regex patterns, otherwise treated as glob. +fn parse_on_finish_of(pattern: &str) -> Option { + let pattern = pattern.trim(); + if let Some(re_str) = pattern.strip_prefix("re:") { + regex::Regex::new(re_str).ok().map(|r| ParsedTrigger::OnFinishOf(Matchable::Regex(r))) + } else { + Glob::new(pattern).ok().and_then(|g| { + let mut builder = GlobSetBuilder::new(); + builder.add(g); + builder.build().ok() + }).map(|gs| ParsedTrigger::OnFinishOf(Matchable::Glob(gs))) + } +} + fn parse_trigger_item_raw(item: &TriggerItemRaw) -> Vec { - match item { - TriggerItemRaw::Simple(s) => match s.as_str() { - "on_success" => vec![ParsedTrigger::OnSuccess], - "on_failure" => vec![ParsedTrigger::OnFailure], - s if s.starts_with("on_finish_of:") => { - let pattern = s.trim_start_matches("on_finish_of:").trim(); - let mut builder = GlobSetBuilder::new(); - builder.add(Glob::new(pattern).expect("Invalid glob pattern")); - let glob_set = builder.build().expect("Failed to compile glob set"); - vec![ParsedTrigger::OnFinishOf(Matchable::Glob(glob_set))] + // Normalize Simple → single-entry Vec, then unify both branches + let entries: Vec<(&str, Vec<&str>)> = match item { + TriggerItemRaw::Simple(s) => { + if let Some(pattern) = s.strip_prefix("on_finish_of:") { + vec![("on_finish_of", vec![pattern.trim()])] + } else { + vec![(s.as_str(), vec![])] } - _ => { - log::warn!("Unknown trigger string: {}", s); - vec![] - } - }, - TriggerItemRaw::Map(map) => { - let mut triggers = Vec::new(); - for (key, patterns) in map { - match key.as_str() { - "on_finish_of" => { - for pattern in patterns { - let mut builder = GlobSetBuilder::new(); - builder.add(Glob::new(pattern).expect("Invalid glob pattern")); - let glob_set = builder.build().expect("Failed to compile glob set"); - triggers.push(ParsedTrigger::OnFinishOf(Matchable::Glob(glob_set))); - } - } - "on_success" => triggers.push(ParsedTrigger::OnSuccess), - "on_failure" => triggers.push(ParsedTrigger::OnFailure), - _ => { - log::warn!("Unknown trigger key: {}", key); + } + TriggerItemRaw::Map(map) => map.iter().map(|(k, v)| + (k.as_str(), v.iter().map(|s| s.as_str()).collect()) + ).collect(), + }; + + let mut triggers = Vec::new(); + for (key, patterns) in entries { + match key { + "on_finish_of" => { + for pattern in patterns { + if let Some(t) = parse_on_finish_of(pattern) { + triggers.push(t); } } } - triggers + "on_success" => triggers.push(ParsedTrigger::OnSuccess), + "on_failure" => triggers.push(ParsedTrigger::OnFailure), + _ => log::warn!("Unknown trigger key: {}", key), } } + triggers } fn merge_templates( @@ -644,48 +628,41 @@ fn merge_templates( ) -> NotificationTemplate { let mut merged = method_template.clone(); if let Some(policy_tmpl) = policy_template - && let Ok(policy_template) = - serde_yaml::from_value::(policy_tmpl.clone()) + && let Ok(policy) = serde_yaml::from_value::(policy_tmpl.clone()) { - if policy_template.schema.is_some() { - merged.schema = policy_template.schema; - } - if policy_template.on_success.is_some() { - merged.on_success = policy_template.on_success; - } - if policy_template.on_failure.is_some() { - merged.on_failure = policy_template.on_failure; - } - if policy_template.on_finish_of.is_some() { - merged.on_finish_of = policy_template.on_finish_of; - } + merged += policy; } merged } -fn extract_overrides(overrides: &serde_yaml::Value) -> (bool, u32, std::time::Duration, f64) { - let fallible = overrides - .get("fallible") - .and_then(|v| v.as_bool()) - .unwrap_or(NOTIFICATION_DEFAULT_FALLIBLE); - let max_lives = overrides - .get("max_lives") - .and_then(|v| v.as_u64()) - .unwrap_or(NOTIFICATION_DEFAULT_MAX_LIVES as u64) as u32; - let timeout_ms = overrides - .get("timeout_ms") - .and_then(|v| v.as_u64()) - .unwrap_or(NOTIFICATION_DEFAULT_TIMEOUT_MS); - let retry_backoff = overrides - .get("retry_backoff") - .and_then(|v| v.as_f64()) - .unwrap_or(NOTIFICATION_DEFAULT_RETRY_BACKOFF); - ( - fallible, - max_lives, - std::time::Duration::from_millis(timeout_ms), - retry_backoff, - ) +struct RuleOverrides { + fallible: bool, + max_lives: u32, + timeout: std::time::Duration, + retry_backoff: f64, +} + +impl RuleOverrides { + fn from_value(v: &serde_yaml::Value) -> Self { + RuleOverrides { + fallible: v.get("fallible").and_then(|v| v.as_bool()).unwrap_or(false), + max_lives: v.get("max_lives").and_then(|v| v.as_u64()).unwrap_or(DEFAULT_NOTIFY_PARAMS.max_lives.value as u64) as u32, + timeout: std::time::Duration::from_millis( + v.get("timeout_ms").and_then(|v| v.as_u64()).unwrap_or(DEFAULT_NOTIFY_PARAMS.timeout_ms.value) + ), + retry_backoff: v.get("retry_backoff").and_then(|v| v.as_f64()).unwrap_or(DEFAULT_NOTIFY_PARAMS.retry_backoff.value), + } + } + + /// Chain: apply `source` atop existing values (source wins when present). + fn apply(self, source: &serde_yaml::Value) -> Self { + RuleOverrides { + fallible: source.get("fallible").and_then(|v| v.as_bool()).unwrap_or(self.fallible), + max_lives: source.get("max_lives").and_then(|v| v.as_u64()).unwrap_or(self.max_lives as u64) as u32, + timeout: source.get("timeout_ms").and_then(|v| v.as_u64()).map(std::time::Duration::from_millis).unwrap_or(self.timeout), + retry_backoff: source.get("retry_backoff").and_then(|v| v.as_f64()).unwrap_or(self.retry_backoff), + } + } } fn parse_rules(config: &NotificationConfig) -> NotificationRulesMap { @@ -694,76 +671,41 @@ fn parse_rules(config: &NotificationConfig) -> NotificationRulesMap { for (priority, group) in &config.groups { let mut ruleset: HashMap> = HashMap::new(); for (method_name, method) in group { - let mut method_rules: Vec = Vec::new(); - if method.policy.is_empty() { - let trigger: Vec = method - .trigger - .iter() - .flat_map(parse_trigger_item_raw) - .collect(); - let template = method.template.clone(); - let (fallible, max_lives, timeout, retry_backoff) = - extract_overrides(&method.config); + let mut method_rules = Vec::new(); + + // Unify: no policy → one empty policy (same merge path) + let empty_policy = crate::finalize::config::NotifyPolicy { + trigger: vec![], + template: serde_yaml::Value::Null, + overrides: serde_yaml::Value::Null, + }; + let policies: &[crate::finalize::config::NotifyPolicy] = if method.policy.is_empty() { + std::slice::from_ref(&empty_policy) + } else { + &method.policy + }; + + for policy in policies { + let trigger: Vec = if policy.trigger.is_empty() { + method.trigger.iter().flat_map(parse_trigger_item_raw).collect() + } else { + policy.trigger.iter().flat_map(parse_trigger_item_raw).collect() + }; + let template = + merge_templates(&method.template, &Some(policy.template.clone())); + let vals = RuleOverrides::from_value(&serde_yaml::Value::Null) + .apply(&method.config) + .apply(&policy.overrides); method_rules.push(NotificationRule { enabled: true, trigger, template, - fallible, - max_lives, - timeout, + fallible: vals.fallible, + max_lives: vals.max_lives, + timeout: vals.timeout, config: method.config.clone(), - retry_backoff, + retry_backoff: vals.retry_backoff, }); - } else { - for policy in &method.policy { - let trigger: Vec = if policy.trigger.is_empty() { - method - .trigger - .iter() - .flat_map(parse_trigger_item_raw) - .collect() - } else { - policy - .trigger - .iter() - .flat_map(parse_trigger_item_raw) - .collect() - }; - let template = - merge_templates(&method.template, &Some(policy.template.clone())); - let method_overrides = extract_overrides(&method.config); - let policy_overrides = extract_overrides(&policy.overrides); - let fallible = if policy.overrides.get("fallible").is_some() { - policy_overrides.0 - } else { - method_overrides.0 - }; - let max_lives = if policy.overrides.get("max_lives").is_some() { - policy_overrides.1 - } else { - method_overrides.1 - }; - let timeout = if policy.overrides.get("timeout_ms").is_some() { - policy_overrides.2 - } else { - method_overrides.2 - }; - let retry_backoff = if policy.overrides.get("retry_backoff").is_some() { - policy_overrides.3 - } else { - method_overrides.3 - }; - method_rules.push(NotificationRule { - enabled: true, - trigger, - template, - fallible, - max_lives, - timeout, - config: method.config.clone(), - retry_backoff, - }); - } } ruleset.insert(method_name.clone(), method_rules); @@ -787,13 +729,17 @@ fn notify_delay(event: &NotificationEvent) -> std::time::Duration { let exponent = event.target.count as i32; if exponent < 0 { log::warn!("Unexpected count: {}", event.target.count); - return std::time::Duration::from_secs_f64(NOTIFICATION_TEMPORARY_MAX_TIME); + return std::time::Duration::from_secs_f64(DEFAULT_NOTIFY_PARAMS.max_retry_delay_seconds.value); } - let factor = if event.retry_backoff > 0.0 { event.retry_backoff } else { NOTIFICATION_TEMPORARY_DELAY_FACTOR as f64 }; + let factor = if event.retry_backoff > 0.0 { + event.retry_backoff + } else { + DEFAULT_NOTIFY_PARAMS.retry_backoff.value + }; let multiplier = factor.powi(exponent); - let delay = multiplier * NOTIFICATION_TEMPORARY_DELAY_TIME; + let delay = multiplier * DEFAULT_NOTIFY_PARAMS.retry_delay_seconds.value; if !delay.is_finite() || delay < 0.0 { - return std::time::Duration::from_secs_f64(NOTIFICATION_TEMPORARY_MAX_TIME); + return std::time::Duration::from_secs_f64(DEFAULT_NOTIFY_PARAMS.max_retry_delay_seconds.value); } std::time::Duration::from_secs_f64(delay) } @@ -834,8 +780,16 @@ mod tests { fn make_tmpl(on_success: &str, on_failure: &str) -> NotificationTemplateDef { NotificationTemplateDef { use_: None, - on_success: if on_success.is_empty() { None } else { Some(on_success.into()) }, - on_failure: if on_failure.is_empty() { None } else { Some(on_failure.into()) }, + on_success: if on_success.is_empty() { + None + } else { + Some(on_success.into()) + }, + on_failure: if on_failure.is_empty() { + None + } else { + Some(on_failure.into()) + }, on_finish_of: None, } } @@ -856,9 +810,9 @@ mod tests { #[test] fn resolve_none_cd_default_found() { let mut ws = HashMap::new(); - ws.insert("default".into(), make_tmpl("OK", "FAIL")); + ws.insert("mail".into(), make_tmpl("OK", "FAIL")); let templates = HashMap::new(); - let result = resolve_template_content(&None, "mail", &templates, false, Some(&ws)); + let result = resolve_template(&None, "mail", &templates, false, Some(&ws)); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "OK"); } @@ -868,7 +822,7 @@ mod tests { let ws: HashMap = HashMap::new(); let mut templates = HashMap::new(); templates.insert("mail.fallback".into(), make_tmpl("FALLBACK", "")); - let result = resolve_template_content(&None, "mail", &templates, false, Some(&ws)); + let result = resolve_template(&None, "mail", &templates, false, Some(&ws)); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "FALLBACK"); } @@ -877,31 +831,29 @@ mod tests { fn resolve_none_cd_no_default_no_fallback_error() { let ws: HashMap = HashMap::new(); let templates = HashMap::new(); - let result = resolve_template_content(&None, "mail", &templates, false, Some(&ws)); + let result = resolve_template(&None, "mail", &templates, false, Some(&ws)); assert!(result.is_err()); - assert!(format!("{}", result.unwrap_err()).contains("no template available")); + assert!(format!("{}", result.unwrap_err()).contains("failed")); } #[test] fn resolve_none_standalone_fallback() { let mut templates = HashMap::new(); templates.insert("mail.fallback".into(), make_tmpl("STANDALONE", "")); - let result = resolve_template_content(&None, "mail", &templates, true, None); + let result = resolve_template(&None, "mail", &templates, true, None); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "STANDALONE"); } #[test] fn resolve_none_standalone_no_fallback_error() { - let result = resolve_template_content(&None, "mail", &HashMap::new(), true, None); + let result = resolve_template(&None, "mail", &HashMap::new(), true, None); assert!(result.is_err()); } #[test] fn resolve_custom() { - let result = resolve_template_content( - &Some("custom".into()), "mail", &HashMap::new(), false, None, - ); + let result = resolve_template(&Some("custom".into()), "mail", &HashMap::new(), false, None); assert!(result.is_ok()); // "custom" returns an empty default — caller reads rule.template.on_* directly } @@ -910,27 +862,27 @@ mod tests { fn resolve_default_standalone_fallback() { let mut templates = HashMap::new(); templates.insert("mail.fallback".into(), make_tmpl("FB", "")); - let result = resolve_template_content( - &Some("default".into()), "mail", &templates, true, None, - ); + let result = resolve_template(&Some("default".into()), "mail", &templates, true, None); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "FB"); } #[test] fn resolve_default_standalone_no_fallback_error() { - let result = resolve_template_content( - &Some("default".into()), "mail", &HashMap::new(), true, None, - ); + let result = resolve_template(&Some("default".into()), "mail", &HashMap::new(), true, None); assert!(result.is_err()); } #[test] fn resolve_default_cd_wsbase() { let mut ws = HashMap::new(); - ws.insert("default".into(), make_tmpl("WS_OK", "")); - let result = resolve_template_content( - &Some("default".into()), "mail", &HashMap::new(), false, Some(&ws), + ws.insert("mail".into(), make_tmpl("WS_OK", "")); + let result = resolve_template( + &Some("default".into()), + "mail", + &HashMap::new(), + false, + Some(&ws), ); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "WS_OK"); @@ -938,8 +890,12 @@ mod tests { #[test] fn resolve_default_cd_no_wsbase_error() { - let result = resolve_template_content( - &Some("default".into()), "mail", &HashMap::new(), false, None, + let result = resolve_template( + &Some("default".into()), + "mail", + &HashMap::new(), + false, + None, ); assert!(result.is_err()); } @@ -948,9 +904,7 @@ mod tests { fn resolve_named_found() { let mut templates = HashMap::new(); templates.insert("sms-tmpl1".into(), make_tmpl("NAMED", "")); - let result = resolve_template_content( - &Some("sms-tmpl1".into()), "mail", &templates, false, None, - ); + let result = resolve_template(&Some("sms-tmpl1".into()), "mail", &templates, false, None); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "NAMED"); } @@ -959,18 +913,14 @@ mod tests { fn resolve_named_not_found_method_match_fallback() { let mut templates = HashMap::new(); templates.insert("mail.fallback".into(), make_tmpl("FB", "")); - let result = resolve_template_content( - &Some("mail".into()), "mail", &templates, false, None, - ); + let result = resolve_template(&Some("mail".into()), "mail", &templates, false, None); assert!(result.is_ok()); assert_eq!(result.unwrap().on_success.unwrap(), "FB"); } #[test] fn resolve_named_not_found_no_match_error() { - let result = resolve_template_content( - &Some("ghost".into()), "mail", &HashMap::new(), false, None, - ); + let result = resolve_template(&Some("ghost".into()), "mail", &HashMap::new(), false, None); assert!(result.is_err()); } @@ -1004,12 +954,15 @@ mod tests { fn flatten_child_overrides_parent() { let mut templates = HashMap::new(); templates.insert("base".into(), make_tmpl("BASE_OK", "BASE_FAIL")); - templates.insert("child".into(), NotificationTemplateDef { - use_: Some("base".into()), - on_success: Some("MY_OK".into()), - on_failure: None, - on_finish_of: None, - }); + templates.insert( + "child".into(), + NotificationTemplateDef { + use_: Some("base".into()), + on_success: Some("MY_OK".into()), + on_failure: None, + on_finish_of: None, + }, + ); let result = flatten_template_refs(&mut templates); assert!(result.is_ok()); assert_eq!(templates["child"].on_success.as_deref(), Some("MY_OK")); @@ -1064,16 +1017,22 @@ mod tests { #[test] fn flatten_external_url_preserved() { let mut templates = HashMap::new(); - templates.insert("ext".into(), NotificationTemplateDef { - use_: Some("git://example.com/repo@v1".into()), - on_success: None, - on_failure: None, - on_finish_of: None, - }); + templates.insert( + "ext".into(), + NotificationTemplateDef { + use_: Some("git://example.com/repo@v1".into()), + on_success: None, + on_failure: None, + on_finish_of: None, + }, + ); let result = flatten_template_refs(&mut templates); assert!(result.is_ok()); // External URLs are NOT flattened — use_ preserved - assert_eq!(templates["ext"].use_.as_deref(), Some("git://example.com/repo@v1")); + assert_eq!( + templates["ext"].use_.as_deref(), + Some("git://example.com/repo@v1") + ); } // ═══════════════════════════════════════════ @@ -1125,7 +1084,11 @@ mod tests { effective_priority: 0, attempted: 0, retry_backoff, - target: NotificationTarget { priority: 0, rules: None, count }, + target: NotificationTarget { + priority: 0, + rules: None, + count, + }, } } diff --git a/workshop-baker/src/notify/interactive.rs b/workshop-baker/src/notify/interactive.rs deleted file mode 100644 index edbb118..0000000 --- a/workshop-baker/src/notify/interactive.rs +++ /dev/null @@ -1,116 +0,0 @@ -use crate::notify::template::render_template; -use crate::notify::error::NotifyError; -use axum::{extract::State, routing::post, Json, Router}; -use serde::{Deserialize, Serialize}; -use std::collections::HashMap; -use std::sync::Arc; -use tokio::sync::oneshot; - -#[derive(Debug, Deserialize)] -struct RenderRequest { - template: String, -} - -#[derive(Debug, Deserialize)] -struct GetRequest { - key: String, -} - -#[derive(Debug, Serialize)] -struct GetResponse { - value: String, -} - -#[derive(Debug, Deserialize)] -struct DoneRequest { - success: bool, -} - -pub struct InteractiveServer { - template_data: Arc>, - shutdown_tx: Option>, - port: u16, -} - -impl InteractiveServer { - pub fn new(template_data: HashMap) -> Self { - Self { - template_data: Arc::new(template_data), - shutdown_tx: None, - port: 0, - } - } - - pub async fn start(&mut self) -> Result { - let (shutdown_tx, shutdown_rx) = oneshot::channel::<()>(); - self.shutdown_tx = Some(shutdown_tx); - - let data = self.template_data.clone(); - - let app = Router::new() - .route("/hbw_plugin_render", post(handle_render)) - .route("/hbw_plugin_get", post(handle_get)) - .route("/hbw_plugin_done", post(handle_done)) - .with_state(data); - - let listener = tokio::net::TcpListener::bind("127.0.0.1:0") - .await - .map_err(|e| NotifyError::TemporaryError(format!("IR server bind failed: {}", e)))?; - let port = listener - .local_addr() - .map_err(|e| NotifyError::TemporaryError(format!("IR server addr error: {}", e)))? - .port(); - self.port = port; - - tokio::spawn(async move { - axum::serve(listener, app) - .with_graceful_shutdown(async move { - let _ = shutdown_rx.await; - }) - .await - .ok(); - }); - - Ok(port) - } - - pub fn port(&self) -> u16 { - self.port - } - - pub async fn shutdown(self) { - if let Some(tx) = self.shutdown_tx { - let _ = tx.send(()); - // Give the server a moment to shut down - tokio::time::sleep(std::time::Duration::from_millis(100)).await; - } - } -} - -async fn handle_render( - State(data): State>>, - Json(req): Json, -) -> Result { - render_template(&req.template, &data) - .map_err(|e| format!("Render error: {}", e)) -} - -async fn handle_get( - State(data): State>>, - Json(req): Json, -) -> Json { - let value = data - .get(&req.key) - .cloned() - .unwrap_or(serde_json::Value::Null); - Json(value) -} - -async fn handle_done(Json(req): Json) -> &'static str { - if req.success { - log::info!("IR plugin reported success"); - } else { - log::warn!("IR plugin reported failure"); - } - "ok" -} diff --git a/workshop-baker/src/notify/types.rs b/workshop-baker/src/notify/types.rs index 9d46993..5100157 100644 --- a/workshop-baker/src/notify/types.rs +++ b/workshop-baker/src/notify/types.rs @@ -91,7 +91,6 @@ impl NotificationRule { #[derive(Debug)] pub enum MethodStatus { Success, - NotAvailable(String), TemporaryError(String), PermanentError(String), Exhausted(String), @@ -164,25 +163,11 @@ impl Ord for NotificationEvent { } } -/// Delay before retrying an unavailable plugin (5 seconds) -pub const NOTIFICATION_NOTREADY_DELAY: std::time::Duration = std::time::Duration::from_secs(5); +use std::sync::LazyLock; +use workshop_baker_params::{NotificationParams, MergeParams}; -/// Exponential backoff base (2.0 = double each retry) -pub const NOTIFICATION_TEMPORARY_DELAY_FACTOR: f64 = 2.0; - -/// Initial backoff delay for temporary errors (5 seconds) -pub const NOTIFICATION_TEMPORARY_DELAY_TIME: f64 = 5.0; - -/// Upper bound for backoff delay (unbounded) -pub const NOTIFICATION_TEMPORARY_MAX_TIME: f64 = f64::MAX; - -/// Cross-group impact factor: remain = max_lives - IMPACT * attempted - count -pub const NOTIFICATION_LIFE_IMPACT_FACTOR: f64 = 0.0; - -pub const NOTIFICATION_DEFAULT_FALLIBLE: bool = false; -pub const NOTIFICATION_DEFAULT_MAX_LIVES: u32 = 3; -pub const NOTIFICATION_DEFAULT_TIMEOUT_MS: u64 = 30000; -pub const NOTIFICATION_DEFAULT_RETRY_BACKOFF: f64 = 2.0; +/// Default parameters used when no external config is provided. +pub static DEFAULT_NOTIFY_PARAMS: LazyLock = LazyLock::new(|| NotificationParams::defaults()); #[derive(Default, Debug, Deserialize, Clone)] pub enum NotificationRenderer {