praxis_protocol/tcp/
mod.rs1use std::{sync::Arc, time::Duration};
7
8use arc_swap::ArcSwap;
9use pingora_core::services::listening::Service;
10use praxis_core::{ProxyError, config::Config};
11use praxis_filter::{FilterPipeline, FilterRegistry};
12use tokio::sync::{Semaphore, watch};
13
14use crate::{ListenerPipelines, Protocol};
15
16pub(crate) mod proxy;
18mod tls_setup;
20
21pub struct PingoraTcp;
32
33#[expect(clippy::too_many_lines, reason = "linear registration with shutdown collection")]
34impl Protocol for PingoraTcp {
35 fn register(
36 self: Box<Self>,
37 server: &mut praxis_core::PingoraServerRuntime,
38 config: &Config,
39 pipelines: &ListenerPipelines,
40 ) -> Result<Vec<watch::Sender<bool>>, ProxyError> {
41 let groups = tls_setup::group_tcp_listeners(config);
42 tls_setup::validate_tcp_group_consistency(&groups)?;
43 #[expect(clippy::expect_used, reason = "empty pipeline is infallible")]
44 let fallback_pipeline = Arc::new(ArcSwap::from_pointee(
45 FilterPipeline::build(&mut [], &FilterRegistry::with_builtins()).expect("empty pipeline is valid"),
46 ));
47
48 let mut cert_watcher_shutdowns = Vec::new();
49
50 for ((upstream_opt, cluster_opt, timeout_ms, max_dur_secs), listeners) in groups {
51 let pipeline = listeners
52 .first()
53 .and_then(|l| pipelines.get(&l.name))
54 .map_or_else(|| Arc::clone(&fallback_pipeline), Arc::clone);
55
56 let session_timeout = timeout_ms.map(Duration::from_millis);
57 let max_duration = max_dur_secs.map(Duration::from_secs);
58 let service_name = match (upstream_opt.as_deref(), cluster_opt.as_deref()) {
59 (Some(addr), _) => format!("tcp-proxy:{addr}"),
60 (_, Some(cluster)) => format!("tcp-proxy:cluster:{cluster}"),
61 _ => "tcp-proxy:filter-routed".to_owned(),
62 };
63 let connection_semaphore = listeners
64 .first()
65 .and_then(|l| l.max_connections)
66 .map(|max| Arc::new(Semaphore::new(max as usize)));
67 let listener_names: std::collections::HashMap<String, ::metrics::SharedString> = listeners
68 .iter()
69 .map(|l| (l.address.clone(), ::metrics::SharedString::from(l.name.clone())))
70 .collect();
71 let default_listener_name = listeners.first().map_or_else(
72 || ::metrics::SharedString::const_str("unknown"),
73 |l| ::metrics::SharedString::from(l.name.clone()),
74 );
75 let app = proxy::PingoraTcpProxy::new(
76 upstream_opt.clone(),
77 cluster_opt.map(Arc::from),
78 pipeline,
79 session_timeout,
80 max_duration,
81 connection_semaphore,
82 config.insecure_options.allow_private_upstreams,
83 listener_names,
84 default_listener_name,
85 );
86 let mut service = Service::new(service_name, app);
87
88 cert_watcher_shutdowns.extend(tls_setup::register_tcp_listeners(
89 &mut service,
90 &listeners,
91 upstream_opt.as_deref(),
92 )?);
93 server.server_mut().add_service(service);
94 }
95
96 Ok(cert_watcher_shutdowns)
97 }
98}