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
21fn upstream_runtime_compatible(old: &UpstreamRuntime, new: &UpstreamConfig) -> bool {
24 *old.chain == new.chain && old.health_config == new.health
25}
26
27pub 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}