1use crate::broker::backend_sdk::{FrameClient, FrameClientError};
31use crate::broker::client::{
32 broker_disabled_by_env, connect_to_backend, BackendConnection, BackendConnectionRoute,
33 BrokerClientError, BrokerDisableEnvError, ConnectBackendRequest,
34};
35use crate::broker::protocol::{Frame, Negotiated};
36
37pub struct BrokerSession {
46 client: FrameClient,
47 route: BackendConnectionRoute,
48 endpoint: String,
49 negotiated: Option<Negotiated>,
50}
51
52impl BrokerSession {
53 pub fn adopt(request: ConnectBackendRequest<'_>) -> Result<Self, AdoptError> {
61 if broker_disabled_by_env()? {
62 return Err(AdoptError::BrokerDisabled);
63 }
64 Ok(Self::from_connection(connect_to_backend(request)?))
65 }
66
67 fn from_connection(connection: BackendConnection) -> Self {
68 Self {
69 client: FrameClient::from_stream(connection.stream),
70 route: connection.route,
71 endpoint: connection.endpoint,
72 negotiated: connection.negotiated,
73 }
74 }
75
76 pub fn route(&self) -> BackendConnectionRoute {
78 self.route
79 }
80
81 pub fn endpoint(&self) -> &str {
83 &self.endpoint
84 }
85
86 pub fn negotiated(&self) -> Option<&Negotiated> {
88 self.negotiated.as_ref()
89 }
90
91 pub fn request(
93 &mut self,
94 payload_protocol: u32,
95 payload: Vec<u8>,
96 ) -> Result<Frame, FrameClientError> {
97 self.client.request(payload_protocol, payload)
98 }
99
100 pub fn client_mut(&mut self) -> &mut FrameClient {
102 &mut self.client
103 }
104
105 pub fn into_client(self) -> FrameClient {
107 self.client
108 }
109
110 pub fn into_backend_io(self) -> Result<OwnedBackendIo, IntoBackendIoError> {
125 let buffered = self.client.buffered_len();
126 if buffered != 0 {
127 return Err(IntoBackendIoError::BufferedResidual { buffered });
128 }
129 OwnedBackendIo::from_local_socket_stream(self.client.into_stream())
130 }
131}
132
133#[derive(Debug)]
141pub struct OwnedBackendIo {
142 #[cfg(unix)]
146 fd: std::os::fd::OwnedFd,
147}
148
149impl OwnedBackendIo {
150 #[cfg(unix)]
151 pub(crate) fn from_local_socket_stream(
152 stream: crate::platform::ipc::Stream,
153 ) -> Result<Self, IntoBackendIoError> {
154 Ok(Self {
155 fd: stream.into_owned_fd(),
156 })
157 }
158
159 #[cfg(windows)]
160 pub(crate) fn from_local_socket_stream(
161 _stream: crate::platform::ipc::Stream,
162 ) -> Result<Self, IntoBackendIoError> {
163 Err(IntoBackendIoError::WindowsUnsupported)
164 }
165
166 #[cfg(unix)]
168 pub fn into_owned_fd(self) -> std::os::fd::OwnedFd {
169 self.fd
170 }
171}
172
173#[cfg(unix)]
174impl std::os::fd::AsFd for OwnedBackendIo {
175 fn as_fd(&self) -> std::os::fd::BorrowedFd<'_> {
176 self.fd.as_fd()
177 }
178}
179
180#[derive(Debug, thiserror::Error)]
183pub enum IntoBackendIoError {
184 #[error(
188 "frame client has {buffered} buffered response byte(s); cannot hand off the raw socket without losing them"
189 )]
190 BufferedResidual {
191 buffered: usize,
193 },
194 #[cfg(feature = "client-async")]
197 #[error("async frame client was poisoned by a prior request panic")]
198 Poisoned,
199 #[cfg(windows)]
202 #[error("into_backend_io() is not yet supported on Windows; the OwnedHandle path is deferred (#720)")]
203 WindowsUnsupported,
204}
205
206#[derive(Debug, thiserror::Error)]
208pub enum AdoptError {
209 #[error("broker disabled via RUNNING_PROCESS_DISABLE=1; use the direct path")]
212 BrokerDisabled,
213 #[error(transparent)]
215 DisableEnv(#[from] BrokerDisableEnvError),
216 #[error(transparent)]
219 Connect(#[from] BrokerClientError),
220 #[cfg(feature = "client-async")]
223 #[error("async adopt worker failed to join: {0}")]
224 AsyncJoin(String),
225}
226
227#[cfg(feature = "client-async")]
234#[derive(Clone, Debug)]
235pub struct OwnedConnectRequest {
236 pub broker_endpoint: String,
238 pub service_name: String,
240 pub wanted_version: String,
242 pub self_version: String,
244 pub cached_backend_endpoint: Option<String>,
246 pub client_version: String,
248 pub client_lib_name: String,
250 pub client_lib_version: String,
252 pub client_keepalive_secs: u64,
254 pub adopt_handed_off_connection: bool,
256 pub handoff_ready_timeout: std::time::Duration,
258}
259
260#[cfg(feature = "client-async")]
261impl OwnedConnectRequest {
262 pub fn new(
264 broker_endpoint: impl Into<String>,
265 service_name: impl Into<String>,
266 wanted_version: impl Into<String>,
267 self_version: impl Into<String>,
268 ) -> Self {
269 Self {
270 broker_endpoint: broker_endpoint.into(),
271 service_name: service_name.into(),
272 wanted_version: wanted_version.into(),
273 self_version: self_version.into(),
274 cached_backend_endpoint: None,
275 client_version: String::new(),
276 client_lib_name: "running-process".to_string(),
277 client_lib_version: env!("CARGO_PKG_VERSION").to_string(),
278 client_keepalive_secs: 0,
279 adopt_handed_off_connection: false,
280 handoff_ready_timeout: crate::broker::client::DEFAULT_HANDOFF_READY_TIMEOUT,
281 }
282 }
283
284 fn as_request(&self) -> ConnectBackendRequest<'_> {
285 ConnectBackendRequest {
286 broker_endpoint: &self.broker_endpoint,
287 service_name: &self.service_name,
288 wanted_version: &self.wanted_version,
289 self_version: &self.self_version,
290 cached_backend_endpoint: self.cached_backend_endpoint.as_deref(),
291 client_version: &self.client_version,
292 client_lib_name: &self.client_lib_name,
293 client_lib_version: &self.client_lib_version,
294 client_keepalive_secs: self.client_keepalive_secs,
295 adopt_handed_off_connection: self.adopt_handed_off_connection,
296 handoff_ready_timeout: self.handoff_ready_timeout,
297 }
298 }
299}
300
301#[cfg(feature = "client-async")]
310pub struct AsyncBrokerSession {
311 client: crate::broker::backend_sdk::AsyncFrameClient,
312 route: BackendConnectionRoute,
313 endpoint: String,
314 negotiated: Option<Negotiated>,
315}
316
317#[cfg(feature = "client-async")]
318impl AsyncBrokerSession {
319 pub async fn adopt(request: OwnedConnectRequest) -> Result<Self, AdoptError> {
322 let joined = tokio::task::spawn_blocking(move || adopt_async_blocking(request))
323 .await
324 .map_err(|err| AdoptError::AsyncJoin(err.to_string()))?;
325 let (route, endpoint, negotiated, client) = joined?;
326 Ok(Self {
327 client: crate::broker::backend_sdk::AsyncFrameClient::from_blocking(client),
328 route,
329 endpoint,
330 negotiated,
331 })
332 }
333
334 pub fn route(&self) -> BackendConnectionRoute {
336 self.route
337 }
338
339 pub fn endpoint(&self) -> &str {
341 &self.endpoint
342 }
343
344 pub fn negotiated(&self) -> Option<&Negotiated> {
346 self.negotiated.as_ref()
347 }
348
349 pub async fn request(
351 &mut self,
352 payload_protocol: u32,
353 payload: Vec<u8>,
354 ) -> Result<Frame, FrameClientError> {
355 self.client.request(payload_protocol, payload).await
356 }
357
358 pub fn into_client(self) -> crate::broker::backend_sdk::AsyncFrameClient {
360 self.client
361 }
362
363 pub fn into_backend_io(self) -> Result<OwnedBackendIo, IntoBackendIoError> {
372 let client = self
373 .client
374 .into_blocking()
375 .ok_or(IntoBackendIoError::Poisoned)?;
376 let buffered = client.buffered_len();
377 if buffered != 0 {
378 return Err(IntoBackendIoError::BufferedResidual { buffered });
379 }
380 OwnedBackendIo::from_local_socket_stream(client.into_stream())
381 }
382}
383
384#[cfg(feature = "client-async")]
385type AdoptedAsync = (
386 BackendConnectionRoute,
387 String,
388 Option<Negotiated>,
389 FrameClient,
390);
391
392#[cfg(feature = "client-async")]
399fn adopt_async_blocking(request: OwnedConnectRequest) -> Result<AdoptedAsync, AdoptError> {
400 if broker_disabled_by_env()? {
401 return Err(AdoptError::BrokerDisabled);
402 }
403
404 #[cfg(feature = "test-seams")]
405 if crate::env_vars::FAKE_BACKEND
406 .os()
407 .is_some_and(|value| !value.is_empty())
408 {
409 return BrokerSession::adopt(request.as_request()).map(|session| {
410 (
411 session.route,
412 session.endpoint,
413 session.negotiated,
414 session.client,
415 )
416 });
417 }
418
419 if request.adopt_handed_off_connection {
420 return BrokerSession::adopt(request.as_request()).map(|session| {
421 (
422 session.route,
423 session.endpoint,
424 session.negotiated,
425 session.client,
426 )
427 });
428 }
429
430 if request.wanted_version == request.self_version {
431 if let Some(endpoint) = request.cached_backend_endpoint.as_deref() {
432 if let Ok(stream) = crate::broker::client::connect_local_socket(endpoint) {
433 return Ok((
434 BackendConnectionRoute::HelloSkip,
435 endpoint.to_owned(),
436 None,
437 FrameClient::from_stream(stream),
438 ));
439 }
440 }
441 }
442
443 let mut hello = request.as_request().hello();
444 hello.request_id = format!("client_v2-{}-{}", request.service_name, std::process::id());
445 let session = crate::broker::client_v2::connect_hello_at_endpoint_with_deadline(
446 request.broker_endpoint,
447 hello,
448 crate::broker::client::broker_client_deadline(),
449 )
450 .map_err(map_explicit_hello_error)?;
451 let negotiated = session.negotiated().clone();
452 let endpoint = negotiated.backend_pipe.clone();
453 let stream = session
454 .connect_backend_ipc()
455 .map_err(map_v2_backend_error)?;
456 Ok((
457 BackendConnectionRoute::BrokerNegotiated,
458 endpoint,
459 Some(negotiated),
460 FrameClient::from_stream(stream),
461 ))
462}
463
464#[cfg(feature = "client-async")]
465fn map_v2_broker_error(error: crate::broker::client_v2::BrokerV2Error) -> AdoptError {
466 use crate::broker::client_v2::BrokerV2Error;
467 let mapped = match error {
468 BrokerV2Error::Dial { source, .. } | BrokerV2Error::Io(source) => {
469 BrokerClientError::BrokerConnect(source)
470 }
471 BrokerV2Error::Framing(source) => BrokerClientError::Framing(source),
472 BrokerV2Error::Decode(source) => BrokerClientError::DecodeHelloReply(source),
473 BrokerV2Error::MissingResult => BrokerClientError::MissingHelloReplyResult,
474 BrokerV2Error::Refused {
475 reason,
476 retry_after_ms,
477 details,
478 } => BrokerClientError::Refused {
479 code: details.code(),
480 reason,
481 retry_after_ms,
482 },
483 other => BrokerClientError::BrokerConnect(std::io::Error::other(other.to_string())),
484 };
485 AdoptError::Connect(mapped)
486}
487
488#[cfg(feature = "client-async")]
489fn map_explicit_hello_error(error: crate::broker::client_v2::ExplicitHelloError) -> AdoptError {
490 use crate::broker::client_v2::ExplicitHelloError;
491 match error {
492 ExplicitHelloError::Broker(error) => map_v2_broker_error(error),
493 ExplicitHelloError::DecodeFrame(source) => {
494 AdoptError::Connect(BrokerClientError::DecodeFrame(source))
495 }
496 ExplicitHelloError::UnexpectedResponseFrame(reason) => {
497 AdoptError::Connect(BrokerClientError::UnexpectedResponseFrame(reason))
498 }
499 }
500}
501
502#[cfg(feature = "client-async")]
503fn map_v2_backend_error(error: crate::broker::client_v2::BackendDialError) -> AdoptError {
504 use crate::broker::client_v2::BackendDialError;
505 let mapped = match error {
506 BackendDialError::EmptyBackendPipe => BrokerClientError::EmptyBackendPipe,
507 BackendDialError::Connect(source) => BrokerClientError::BackendConnect(source),
508 BackendDialError::IntoBackendIo(source) => {
509 BrokerClientError::BackendConnect(std::io::Error::other(source.to_string()))
510 }
511 };
512 AdoptError::Connect(mapped)
513}