running_process/broker/server/
serve.rs1use std::num::NonZeroUsize;
10use std::path::PathBuf;
11use std::sync::Mutex;
12use std::time::{Duration, Instant};
13
14use crate::broker::backend_handle::{BackendHandle, BackendHandleError, DaemonProcess};
15use crate::broker::backend_lifecycle::identity::IdentityError;
16use crate::broker::lifecycle::sid::SidError;
17use crate::broker::protocol::{Endpoint, ServiceDefinition};
18
19use super::admin::AdminSnapshot;
20use super::backend_launcher::{BackendLauncher, CommandBackendLauncher};
21use super::backend_registry::BackendRegistry;
22use super::combined_service_def_loader::CombinedServiceDefinitionLoader;
23use super::connection::{BrokerConnectionError, PeerCredentialPolicy};
24use super::control_socket::{
25 serve_control_socket_connections_with_limit_policy_post_hello_opaque,
26 serve_launch_control_socket_connections_concurrently, ControlSocketConnectionLimit,
27 ControlSocketError,
28};
29use super::fd_pressure::FdPressureGuard;
30use super::handoff_serve::{try_complete_negotiated_handoff_opaque, ServeHandoffContext};
31use super::hello_handler::{HelloHandler, HelloHandlerError};
32use super::hello_router::HelloRouter;
33use super::instance::{BrokerInstanceError, BrokerInstanceKey};
34use super::service_def_loader::{service_definition_dir, ServiceDefinitionError};
35use super::spawn_coordinator::SpawnCoordinator;
36use super::version_allow_list::{check_version_allowed, VersionPolicyBlock};
37
38#[derive(Clone, Debug)]
40pub struct BrokerServeConfig {
41 pub socket_path: String,
43 pub service_name: String,
45 pub service_version: String,
47 pub backend_endpoint: String,
49 pub service_definition_dir: PathBuf,
51 pub max_connections: Option<NonZeroUsize>,
53 pub handoff_endpoint: Option<String>,
58 pub endpoint_probe_timeout: Option<Duration>,
67}
68
69#[derive(Clone, Debug)]
71pub struct BrokerLaunchServeConfig {
72 pub socket_path: String,
74 pub service_definition_dir: PathBuf,
76 pub max_connections: Option<NonZeroUsize>,
78}
79
80impl BrokerServeConfig {
81 pub fn new(
83 socket_path: impl Into<String>,
84 service_name: impl Into<String>,
85 service_version: impl Into<String>,
86 backend_endpoint: impl Into<String>,
87 max_connections: usize,
88 ) -> Result<Self, BrokerServeError> {
89 Ok(Self {
90 socket_path: socket_path.into(),
91 service_name: service_name.into(),
92 service_version: service_version.into(),
93 backend_endpoint: backend_endpoint.into(),
94 service_definition_dir: service_definition_dir(),
95 max_connections: Some(
96 NonZeroUsize::new(max_connections)
97 .ok_or(BrokerServeError::InvalidMaxConnections)?,
98 ),
99 handoff_endpoint: None,
100 endpoint_probe_timeout: None,
101 })
102 }
103
104 pub fn unbounded(
107 socket_path: impl Into<String>,
108 service_name: impl Into<String>,
109 service_version: impl Into<String>,
110 backend_endpoint: impl Into<String>,
111 ) -> Self {
112 Self {
113 socket_path: socket_path.into(),
114 service_name: service_name.into(),
115 service_version: service_version.into(),
116 backend_endpoint: backend_endpoint.into(),
117 service_definition_dir: service_definition_dir(),
118 max_connections: None,
119 handoff_endpoint: None,
120 endpoint_probe_timeout: None,
121 }
122 }
123
124 pub fn with_service_definition_dir(mut self, root: impl Into<PathBuf>) -> Self {
126 self.service_definition_dir = root.into();
127 self
128 }
129
130 pub fn with_handoff_endpoint(mut self, endpoint: impl Into<String>) -> Self {
133 self.handoff_endpoint = Some(endpoint.into());
134 self
135 }
136
137 pub fn with_endpoint_probe_timeout(mut self, timeout: Duration) -> Self {
141 self.endpoint_probe_timeout = Some(timeout);
142 self
143 }
144
145 pub fn connection_limit(&self) -> ControlSocketConnectionLimit {
147 self.max_connections.map_or(
148 ControlSocketConnectionLimit::Unbounded,
149 ControlSocketConnectionLimit::Bounded,
150 )
151 }
152}
153
154impl BrokerLaunchServeConfig {
155 pub fn new(
158 socket_path: impl Into<String>,
159 max_connections: usize,
160 ) -> Result<Self, BrokerServeError> {
161 Ok(Self {
162 socket_path: socket_path.into(),
163 service_definition_dir: service_definition_dir(),
164 max_connections: Some(
165 NonZeroUsize::new(max_connections)
166 .ok_or(BrokerServeError::InvalidMaxConnections)?,
167 ),
168 })
169 }
170
171 pub fn unbounded(socket_path: impl Into<String>) -> Self {
174 Self {
175 socket_path: socket_path.into(),
176 service_definition_dir: service_definition_dir(),
177 max_connections: None,
178 }
179 }
180
181 pub fn with_service_definition_dir(mut self, root: impl Into<PathBuf>) -> Self {
183 self.service_definition_dir = root.into();
184 self
185 }
186
187 pub fn connection_limit(&self) -> ControlSocketConnectionLimit {
189 self.max_connections.map_or(
190 ControlSocketConnectionLimit::Unbounded,
191 ControlSocketConnectionLimit::Bounded,
192 )
193 }
194}
195
196pub fn serve_registered_backend(config: BrokerServeConfig) -> Result<(), BrokerServeError> {
198 let RegisteredServeBackend {
199 loader,
200 registry,
201 instance,
202 ..
203 } = build_registered_backend(&config)?;
204 let registry = Mutex::new(registry);
205 let router = HelloRouter::with_lifecycle_monitor(&loader, ®istry);
206 let peer_policy =
207 PeerCredentialPolicy::current_user().ok_or(BrokerServeError::PeerPolicyUnavailable)?;
208 let started_at = Instant::now();
209 let fd_guard = FdPressureGuard::default();
210 let snapshot_provider = || {
211 let registry = registry
212 .lock()
213 .unwrap_or_else(|poisoned| poisoned.into_inner());
214 let demoted = fd_guard.is_demoted();
215 AdminSnapshot::from_registry(
216 instance.id(),
217 started_at.elapsed(),
218 !demoted,
219 0,
220 ®istry,
221 &[],
222 )
223 .with_fd_pressure_demoted(demoted)
224 };
225 serve_control_socket_connections_with_limit_policy_post_hello_opaque(
226 &config.socket_path,
227 &router,
228 snapshot_provider,
229 config.connection_limit(),
230 &peer_policy,
231 |mut stream, reply| {
232 let Some(handoff_endpoint) = config.handoff_endpoint.as_deref() else {
234 return;
235 };
236 let ctx = ServeHandoffContext {
237 handoff_endpoint,
238 service_name: &config.service_name,
239 service_version: &config.service_version,
240 instance: &instance,
241 registry: ®istry,
242 };
243 let _must_relinquish = try_complete_negotiated_handoff_opaque(&ctx, &mut stream, reply);
244 },
245 &fd_guard,
246 )?;
247 Ok(())
248}
249
250pub fn serve_launching_backends(config: BrokerLaunchServeConfig) -> Result<(), BrokerServeError> {
253 let launcher = CommandBackendLauncher::for_current_user()?;
254 serve_launching_backends_with_launcher(config, &launcher)
255}
256
257pub fn serve_launching_backends_with_launcher(
259 config: BrokerLaunchServeConfig,
260 launcher: &dyn BackendLauncher,
261) -> Result<(), BrokerServeError> {
262 let loader = CombinedServiceDefinitionLoader::new(&config.service_definition_dir);
263 let registry = Mutex::new(BackendRegistry::new());
264 let spawn_coordinator = Mutex::new(SpawnCoordinator::new());
265 let router = HelloRouter::with_lifecycle_monitor(&loader, ®istry)
266 .with_spawn_coordinator(&spawn_coordinator)
267 .with_backend_launcher(launcher);
268 let peer_policy =
269 PeerCredentialPolicy::current_user().ok_or(BrokerServeError::PeerPolicyUnavailable)?;
270 let started_at = Instant::now();
271 let fd_guard = FdPressureGuard::default();
272 let snapshot_provider = || {
273 let registry = registry
274 .lock()
275 .unwrap_or_else(|poisoned| poisoned.into_inner());
276 let demoted = fd_guard.is_demoted();
277 AdminSnapshot::from_registry("launch", started_at.elapsed(), !demoted, 0, ®istry, &[])
278 .with_fd_pressure_demoted(demoted)
279 };
280 serve_launch_control_socket_connections_concurrently(
281 &config.socket_path,
282 &router,
283 snapshot_provider,
284 config.connection_limit(),
285 &peer_policy,
286 &fd_guard,
287 )?;
288 Ok(())
289}
290
291pub fn build_hello_handler(config: &BrokerServeConfig) -> Result<HelloHandler, BrokerServeError> {
293 let registered = build_registered_backend(config)?;
294 let backend = registered
295 .registry
296 .registered_backend_for_any_build(
297 ®istered.instance,
298 ®istered.service_definition,
299 &config.service_version,
300 )
301 .ok_or(BrokerServeError::RegisteredBackendMissing)?;
302
303 Ok(HelloHandler::new().with_backend(backend)?)
304}
305
306struct RegisteredServeBackend {
307 loader: CombinedServiceDefinitionLoader,
308 registry: BackendRegistry,
309 instance: BrokerInstanceKey,
310 service_definition: ServiceDefinition,
311}
312
313fn build_registered_backend(
314 config: &BrokerServeConfig,
315) -> Result<RegisteredServeBackend, BrokerServeError> {
316 if config.backend_endpoint.is_empty() {
317 return Err(BrokerServeError::EmptyBackendEndpoint);
318 }
319
320 let loader = CombinedServiceDefinitionLoader::new(&config.service_definition_dir);
321 let service_definition = loader.lookup_or_reload(&config.service_name)?;
322 check_version_allowed(&config.service_version, &service_definition)
323 .map_err(BrokerServeError::VersionPolicy)?;
324
325 let instance = BrokerInstanceKey::from_service_definition(&service_definition)?;
326 let endpoint = Endpoint {
327 namespace_id: instance.id(),
328 path: config.backend_endpoint.clone(),
329 };
330 let daemon = DaemonProcess::current_process(endpoint.clone(), Some(30))?;
331 let handle = BackendHandle::probe_with_service_and_timeout(
332 config.service_name.clone(),
333 config.service_version.clone(),
334 &endpoint,
335 &daemon,
336 config
337 .endpoint_probe_timeout
338 .unwrap_or(crate::broker::backend_lifecycle::probe::DEFAULT_ENDPOINT_PROBE_TIMEOUT),
339 )?;
340
341 let mut registry = BackendRegistry::new();
342 registry.insert(instance.clone(), handle);
343
344 Ok(RegisteredServeBackend {
345 loader,
346 registry,
347 instance,
348 service_definition,
349 })
350}
351
352#[derive(Debug, thiserror::Error)]
354pub enum BrokerServeError {
355 #[error("max_connections must be greater than zero")]
357 InvalidMaxConnections,
358 #[error("backend endpoint must not be empty")]
360 EmptyBackendEndpoint,
361 #[error(transparent)]
363 ServiceDefinition(#[from] ServiceDefinitionError),
364 #[error(transparent)]
366 BrokerInstance(#[from] BrokerInstanceError),
367 #[error(transparent)]
369 Identity(#[from] IdentityError),
370 #[error(transparent)]
372 Sid(#[from] SidError),
373 #[error("configured service version is blocked by service-definition policy: {0:?}")]
375 VersionPolicy(VersionPolicyBlock),
376 #[error(transparent)]
378 BackendHandle(#[from] BackendHandleError),
379 #[error("registered backend was missing after registry insert")]
381 RegisteredBackendMissing,
382 #[error(transparent)]
384 HelloHandler(#[from] HelloHandlerError),
385 #[error("current-user peer credential policy is unavailable")]
387 PeerPolicyUnavailable,
388 #[error(transparent)]
390 Connection(#[from] BrokerConnectionError),
391 #[error(transparent)]
393 ControlSocket(#[from] ControlSocketError),
394}