Skip to main content

eggress_runtime/
snapshot.rs

1use std::collections::HashMap;
2use std::sync::Arc;
3
4use eggress_config::compile::{
5    AdminConfig, GroupFallback, ListenerConfig, RuntimeConfig, UpstreamConfig,
6};
7use eggress_routing::upstream::{UpstreamGroup, UpstreamRuntime};
8use eggress_routing::{RouteActionSpec, Router};
9
10pub struct CompiledRuntimeSnapshot {
11    pub generation: u64,
12    pub upstreams: HashMap<String, Arc<UpstreamRuntime>>,
13    pub router: Arc<Router>,
14    pub timeouts: eggress_config::compile::TimeoutConfig,
15    pub listeners: Vec<ListenerConfig>,
16    pub admin: Option<AdminConfig>,
17    pub reverse_servers: Vec<eggress_config::compile::CompiledReverseServerConfig>,
18    pub reverse_clients: Vec<eggress_config::compile::CompiledReverseClientConfig>,
19}
20
21/// Check whether an existing `UpstreamRuntime` is compatible with a new config,
22/// meaning its chain specification hasn't changed and we can reuse the Arc.
23fn upstream_runtime_compatible(old: &UpstreamRuntime, new: &UpstreamConfig) -> bool {
24    *old.chain == new.chain && old.health_config == new.health
25}
26
27/// Build a `CompiledRuntimeSnapshot` from a `RuntimeConfig`.
28///
29/// Upstream runtimes are created first and shared with groups/router so that
30/// the same `Arc<UpstreamRuntime>` objects are used for health probing and routing.
31pub fn compile_runtime_snapshot(
32    rt: &RuntimeConfig,
33    previous: Option<&CompiledRuntimeSnapshot>,
34) -> Result<CompiledRuntimeSnapshot, Box<dyn std::error::Error + Send + Sync>> {
35    let empty_map = HashMap::new();
36    let previous_upstreams = previous.map(|p| &p.upstreams).unwrap_or(&empty_map);
37
38    let mut upstreams: HashMap<String, Arc<UpstreamRuntime>> = HashMap::new();
39
40    for u in &rt.upstreams {
41        let runtime = if let Some(existing) = previous_upstreams.get(&u.id) {
42            if upstream_runtime_compatible(existing, u) {
43                existing.clone()
44            } else {
45                build_one_upstream_runtime(u)
46            }
47        } else {
48            build_one_upstream_runtime(u)
49        };
50        upstreams.insert(u.id.clone(), runtime);
51    }
52
53    let group_ids: std::collections::HashSet<_> = rt.groups.iter().map(|g| g.id.clone()).collect();
54
55    let mut groups = Vec::new();
56    for g in &rt.groups {
57        let mut members = Vec::new();
58        for m in &g.members {
59            let member = upstreams
60                .get(m)
61                .ok_or_else(|| format!("group '{}' references unknown upstream '{}'", g.id, m))?;
62            members.push(member.clone());
63        }
64        if members.is_empty() {
65            return Err(format!("group '{}' has no valid members", g.id).into());
66        }
67
68        let fallback = match g.fallback {
69            GroupFallback::Reject => eggress_routing::upstream::GroupFallback::Reject,
70            GroupFallback::Direct => eggress_routing::upstream::GroupFallback::Direct,
71            GroupFallback::UseUnhealthy => eggress_routing::upstream::GroupFallback::UseUnhealthy,
72        };
73
74        groups.push((
75            g.id.clone(),
76            UpstreamGroup::new(g.id.clone(), g.scheduler, Arc::from(members), fallback),
77        ));
78    }
79
80    let mut rules = Vec::new();
81    for r in &rt.rules {
82        let action = match &r.action {
83            RouteActionSpec::Direct => RouteActionSpec::Direct,
84            RouteActionSpec::UpstreamGroup(gid) => {
85                if !group_ids.contains(gid) {
86                    return Err(
87                        format!("rule '{}' references unknown group '{}'", r.id, gid).into(),
88                    );
89                }
90                RouteActionSpec::UpstreamGroup(gid.clone())
91            }
92            RouteActionSpec::Reject(reason) => RouteActionSpec::Reject(reason.clone()),
93        };
94        rules.push(eggress_routing::CompiledRule {
95            id: r.id.clone(),
96            matcher: r.matcher.clone(),
97            action,
98        });
99    }
100
101    let router = Router::with_groups(rules, rt.default_action.clone(), groups);
102    let gen = previous.map(|p| p.generation + 1).unwrap_or(0);
103
104    Ok(CompiledRuntimeSnapshot {
105        generation: gen,
106        upstreams,
107        router: Arc::new(router),
108        timeouts: rt.timeouts.clone(),
109        listeners: rt.listeners.clone(),
110        admin: rt.admin.clone(),
111        reverse_servers: rt.reverse_servers.clone(),
112        reverse_clients: rt.reverse_clients.clone(),
113    })
114}
115
116fn build_one_upstream_runtime(u: &UpstreamConfig) -> Arc<UpstreamRuntime> {
117    let id = eggress_core::UpstreamId::new(u.id.clone());
118    let mut runtime =
119        UpstreamRuntime::new(id, u.chain.clone()).with_health_config(u.health.clone());
120
121    if let Some(first_hop) = u.chain.hops.first() {
122        match format!("{}:{}", first_hop.endpoint.host, first_hop.endpoint.port)
123            .parse::<std::net::SocketAddr>()
124        {
125            Ok(addr) => {
126                runtime =
127                    runtime.with_health_probe(eggress_routing::health::HealthProbe::TcpConnect {
128                        target: addr,
129                        timeout: u.health.timeout,
130                    });
131            }
132            Err(_) => {
133                tracing::warn!(
134                    upstream = %u.id,
135                    endpoint = %first_hop.endpoint.host,
136                    "skipping health probe for hostname endpoint; upstream stays at its initial health state"
137                );
138            }
139        }
140    }
141
142    Arc::new(runtime)
143}
144
145#[cfg(test)]
146mod tests {
147    use super::*;
148    use eggress_config::compile::{
149        GroupFallback, ProcessConfig, RuntimeConfig, TimeoutConfig, UpstreamConfig,
150    };
151    use eggress_routing::scheduler::SchedulerKind;
152    use eggress_routing::UpstreamGroupId;
153    use eggress_uri::ProxyChainSpec;
154    use std::time::Duration;
155
156    fn default_health() -> eggress_routing::health::HealthConfig {
157        eggress_routing::health::HealthConfig::default()
158    }
159
160    fn empty_config() -> RuntimeConfig {
161        RuntimeConfig {
162            process: ProcessConfig::default(),
163            timeouts: TimeoutConfig::default(),
164            listeners: vec![],
165            upstreams: vec![],
166            groups: vec![],
167            rules: vec![],
168            default_action: RouteActionSpec::Direct,
169            admin: None,
170            reverse_servers: vec![],
171            reverse_clients: vec![],
172        }
173    }
174
175    #[test]
176    fn snapshot_empty_config() {
177        let snap = compile_runtime_snapshot(&empty_config(), None).unwrap();
178        assert_eq!(snap.generation, 0);
179        assert!(snap.upstreams.is_empty());
180        assert!(snap.router.rules().is_empty());
181    }
182
183    #[test]
184    fn snapshot_single_upstream() {
185        let mut cfg = empty_config();
186        cfg.upstreams = vec![UpstreamConfig {
187            id: "proxy1".to_string(),
188            chain: ProxyChainSpec { hops: vec![] },
189            health: default_health(),
190            h2: None,
191        }];
192        let snap = compile_runtime_snapshot(&cfg, None).unwrap();
193        assert_eq!(snap.upstreams.len(), 1);
194        assert!(snap.upstreams.contains_key("proxy1"));
195    }
196
197    #[test]
198    fn snapshot_group_uses_shared_upstream_arc() {
199        let mut cfg = empty_config();
200        cfg.upstreams = vec![UpstreamConfig {
201            id: "proxy1".to_string(),
202            chain: ProxyChainSpec { hops: vec![] },
203            health: default_health(),
204            h2: None,
205        }];
206        cfg.groups = vec![eggress_config::compile::UpstreamGroupConfig {
207            id: UpstreamGroupId(Arc::from("main")),
208            scheduler: SchedulerKind::RoundRobin,
209            members: vec!["proxy1".to_string()],
210            fallback: GroupFallback::Reject,
211        }];
212        let snap = compile_runtime_snapshot(&cfg, None).unwrap();
213        let upstream_arc = snap.upstreams.get("proxy1").unwrap();
214        let group = snap
215            .router
216            .groups()
217            .get(&UpstreamGroupId(Arc::from("main")))
218            .unwrap();
219        let group_member = &group.members[0];
220        assert!(Arc::ptr_eq(upstream_arc, group_member));
221    }
222
223    #[test]
224    fn unchanged_upstream_retains_arc_identity_after_reload() {
225        let mut cfg = empty_config();
226        cfg.upstreams = vec![UpstreamConfig {
227            id: "proxy1".to_string(),
228            chain: ProxyChainSpec { hops: vec![] },
229            health: default_health(),
230            h2: None,
231        }];
232        let snap1 = compile_runtime_snapshot(&cfg, None).unwrap();
233        let original_arc = snap1.upstreams.get("proxy1").unwrap().clone();
234
235        let snap2 = compile_runtime_snapshot(&cfg, Some(&snap1)).unwrap();
236        let reused_arc = snap2.upstreams.get("proxy1").unwrap().clone();
237
238        assert!(Arc::ptr_eq(&original_arc, &reused_arc));
239    }
240
241    #[test]
242    fn changed_upstream_gets_fresh_arc() {
243        let mut cfg1 = empty_config();
244        cfg1.upstreams = vec![UpstreamConfig {
245            id: "proxy1".to_string(),
246            chain: ProxyChainSpec { hops: vec![] },
247            health: default_health(),
248            h2: None,
249        }];
250        let snap1 = compile_runtime_snapshot(&cfg1, None).unwrap();
251        let original_arc = snap1.upstreams.get("proxy1").unwrap().clone();
252
253        let mut cfg2 = empty_config();
254        cfg2.upstreams = vec![UpstreamConfig {
255            id: "proxy1".to_string(),
256            chain: ProxyChainSpec {
257                hops: vec![eggress_uri::ProxyHopSpec {
258                    protocols: vec![eggress_uri::ProtocolSpec::Socks5],
259                    endpoint: eggress_uri::EndpointSpec {
260                        host: "newhost".to_string(),
261                        port: 1080,
262                    },
263                    credentials: None,
264                    rule: None,
265                    local_bind: None,
266                    tls: false,
267                    server_name: None,
268                    insecure: false,
269                    plugins: Vec::new(),
270                    auth_prefix: None,
271                }],
272            },
273            health: default_health(),
274            h2: None,
275        }];
276        let snap2 = compile_runtime_snapshot(&cfg2, Some(&snap1)).unwrap();
277        let new_arc = snap2.upstreams.get("proxy1").unwrap().clone();
278
279        assert!(!Arc::ptr_eq(&original_arc, &new_arc));
280    }
281
282    #[test]
283    fn no_duplicate_upstream_runtime_objects_for_one_id() {
284        let mut cfg = empty_config();
285        cfg.upstreams = vec![UpstreamConfig {
286            id: "proxy1".to_string(),
287            chain: ProxyChainSpec { hops: vec![] },
288            health: default_health(),
289            h2: None,
290        }];
291        cfg.groups = vec![eggress_config::compile::UpstreamGroupConfig {
292            id: UpstreamGroupId(Arc::from("main")),
293            scheduler: SchedulerKind::RoundRobin,
294            members: vec!["proxy1".to_string()],
295            fallback: GroupFallback::Reject,
296        }];
297        let snap = compile_runtime_snapshot(&cfg, None).unwrap();
298        let upstream_arc = snap.upstreams.get("proxy1").unwrap();
299        let group = snap
300            .router
301            .groups()
302            .get(&UpstreamGroupId(Arc::from("main")))
303            .unwrap();
304        let group_member = &group.members[0];
305
306        assert!(Arc::ptr_eq(upstream_arc, group_member));
307    }
308
309    #[test]
310    fn generation_increments_on_reload() {
311        let cfg = empty_config();
312        let snap1 = compile_runtime_snapshot(&cfg, None).unwrap();
313        assert_eq!(snap1.generation, 0);
314        let snap2 = compile_runtime_snapshot(&cfg, Some(&snap1)).unwrap();
315        assert_eq!(snap2.generation, 1);
316        let snap3 = compile_runtime_snapshot(&cfg, Some(&snap2)).unwrap();
317        assert_eq!(snap3.generation, 2);
318    }
319
320    #[test]
321    fn group_references_unknown_upstream() {
322        let mut cfg = empty_config();
323        cfg.groups = vec![eggress_config::compile::UpstreamGroupConfig {
324            id: UpstreamGroupId(Arc::from("main")),
325            scheduler: SchedulerKind::RoundRobin,
326            members: vec!["nonexistent".to_string()],
327            fallback: GroupFallback::Reject,
328        }];
329        let result = compile_runtime_snapshot(&cfg, None);
330        assert!(result.is_err());
331        let err_msg = result.err().unwrap().to_string();
332        assert!(err_msg.contains("nonexistent"));
333    }
334
335    #[test]
336    fn rule_references_unknown_group() {
337        let mut cfg = empty_config();
338        cfg.rules = vec![eggress_routing::CompiledRule {
339            id: eggress_routing::RuleId(Arc::from("r1")),
340            matcher: eggress_routing::MatchExpr::Any,
341            action: RouteActionSpec::UpstreamGroup(UpstreamGroupId(Arc::from("missing"))),
342        }];
343        let result = compile_runtime_snapshot(&cfg, None);
344        assert!(result.is_err());
345        let err_msg = result.err().unwrap().to_string();
346        assert!(err_msg.contains("missing"));
347    }
348
349    #[test]
350    fn multiple_upstreams_all_shared() {
351        let mut cfg = empty_config();
352        cfg.upstreams = vec![
353            UpstreamConfig {
354                id: "p1".to_string(),
355                chain: ProxyChainSpec { hops: vec![] },
356                health: default_health(),
357                h2: None,
358            },
359            UpstreamConfig {
360                id: "p2".to_string(),
361                chain: ProxyChainSpec { hops: vec![] },
362                health: default_health(),
363                h2: None,
364            },
365        ];
366        cfg.groups = vec![eggress_config::compile::UpstreamGroupConfig {
367            id: UpstreamGroupId(Arc::from("grp")),
368            scheduler: SchedulerKind::RoundRobin,
369            members: vec!["p1".to_string(), "p2".to_string()],
370            fallback: GroupFallback::Reject,
371        }];
372        let snap = compile_runtime_snapshot(&cfg, None).unwrap();
373        let group = snap
374            .router
375            .groups()
376            .get(&UpstreamGroupId(Arc::from("grp")))
377            .unwrap();
378        assert!(Arc::ptr_eq(
379            snap.upstreams.get("p1").unwrap(),
380            &group.members[0]
381        ));
382        assert!(Arc::ptr_eq(
383            snap.upstreams.get("p2").unwrap(),
384            &group.members[1]
385        ));
386    }
387
388    #[test]
389    fn changed_health_config_gets_fresh_arc() {
390        let mut cfg1 = empty_config();
391        cfg1.upstreams = vec![UpstreamConfig {
392            id: "proxy1".to_string(),
393            chain: ProxyChainSpec { hops: vec![] },
394            health: eggress_routing::health::HealthConfig {
395                failures_to_unhealthy: 3,
396                ..default_health()
397            },
398            h2: None,
399        }];
400        let snap1 = compile_runtime_snapshot(&cfg1, None).unwrap();
401        let original_arc = snap1.upstreams.get("proxy1").unwrap().clone();
402
403        let mut cfg2 = empty_config();
404        cfg2.upstreams = vec![UpstreamConfig {
405            id: "proxy1".to_string(),
406            chain: ProxyChainSpec { hops: vec![] },
407            health: eggress_routing::health::HealthConfig {
408                failures_to_unhealthy: 1,
409                ..default_health()
410            },
411            h2: None,
412        }];
413        let snap2 = compile_runtime_snapshot(&cfg2, Some(&snap1)).unwrap();
414        let new_arc = snap2.upstreams.get("proxy1").unwrap().clone();
415
416        assert!(
417            !Arc::ptr_eq(&original_arc, &new_arc),
418            "changing health config should produce a fresh upstream runtime ARC"
419        );
420    }
421
422    #[test]
423    fn unchanged_health_config_retains_arc_identity() {
424        let mut cfg1 = empty_config();
425        cfg1.upstreams = vec![UpstreamConfig {
426            id: "proxy1".to_string(),
427            chain: ProxyChainSpec { hops: vec![] },
428            health: eggress_routing::health::HealthConfig {
429                failures_to_unhealthy: 2,
430                successes_to_healthy: 1,
431                interval: Duration::from_secs(10),
432                timeout: Duration::from_secs(2),
433                initial_state: eggress_routing::health::HealthState::Unknown,
434            },
435            h2: None,
436        }];
437        let snap1 = compile_runtime_snapshot(&cfg1, None).unwrap();
438        let original_arc = snap1.upstreams.get("proxy1").unwrap().clone();
439
440        let mut cfg2 = empty_config();
441        cfg2.upstreams = vec![UpstreamConfig {
442            id: "proxy1".to_string(),
443            chain: ProxyChainSpec { hops: vec![] },
444            health: eggress_routing::health::HealthConfig {
445                failures_to_unhealthy: 2,
446                successes_to_healthy: 1,
447                interval: Duration::from_secs(10),
448                timeout: Duration::from_secs(2),
449                initial_state: eggress_routing::health::HealthState::Unknown,
450            },
451            h2: None,
452        }];
453        let snap2 = compile_runtime_snapshot(&cfg2, Some(&snap1)).unwrap();
454        let reused_arc = snap2.upstreams.get("proxy1").unwrap().clone();
455
456        assert!(
457            Arc::ptr_eq(&original_arc, &reused_arc),
458            "identical health config should retain the upstream runtime ARC"
459        );
460    }
461
462    #[test]
463    fn partial_upstream_change_preserves_others() {
464        let mut cfg1 = empty_config();
465        cfg1.upstreams = vec![
466            UpstreamConfig {
467                id: "p1".to_string(),
468                chain: ProxyChainSpec { hops: vec![] },
469                health: default_health(),
470                h2: None,
471            },
472            UpstreamConfig {
473                id: "p2".to_string(),
474                chain: ProxyChainSpec { hops: vec![] },
475                health: default_health(),
476                h2: None,
477            },
478        ];
479        let snap1 = compile_runtime_snapshot(&cfg1, None).unwrap();
480        let p1_original = snap1.upstreams.get("p1").unwrap().clone();
481        let p2_original = snap1.upstreams.get("p2").unwrap().clone();
482
483        let mut cfg2 = empty_config();
484        cfg2.upstreams = vec![
485            UpstreamConfig {
486                id: "p1".to_string(),
487                chain: ProxyChainSpec { hops: vec![] },
488                health: default_health(),
489                h2: None,
490            },
491            UpstreamConfig {
492                id: "p2".to_string(),
493                chain: ProxyChainSpec {
494                    hops: vec![eggress_uri::ProxyHopSpec {
495                        protocols: vec![eggress_uri::ProtocolSpec::Http],
496                        endpoint: eggress_uri::EndpointSpec {
497                            host: "changed".to_string(),
498                            port: 8080,
499                        },
500                        credentials: None,
501                        rule: None,
502                        local_bind: None,
503                        tls: false,
504                        server_name: None,
505                        insecure: false,
506                        plugins: Vec::new(),
507                        auth_prefix: None,
508                    }],
509                },
510                health: default_health(),
511                h2: None,
512            },
513        ];
514        let snap2 = compile_runtime_snapshot(&cfg2, Some(&snap1)).unwrap();
515
516        assert!(Arc::ptr_eq(
517            &p1_original,
518            snap2.upstreams.get("p1").unwrap()
519        ));
520        assert!(!Arc::ptr_eq(
521            &p2_original,
522            snap2.upstreams.get("p2").unwrap()
523        ));
524    }
525}