running_process/broker/backend_lib/
wire.rs1use std::io::{Read, Write};
20use std::time::Instant;
21
22use prost::Message;
23
24use crate::broker::backend_lib::accept_handed_off::{
25 accept_handed_off, HandedOffPayload, HandoffAcceptance,
26};
27use crate::broker::protocol::{
28 read_frame, write_frame, Frame, FrameKind, FramingError, HandoffAck, HandoffOffer,
29};
30use crate::broker::server::deadline_stream::with_nonblocking_deadline;
31use crate::broker::server::handoff::wire::{handoff_ack_frame, validate_handoff_frame};
32use crate::broker::server::{HandoffToken, HandoffTokenStore};
33
34#[derive(Debug, thiserror::Error)]
36pub enum BackendHandoffWireError {
37 #[error(transparent)]
39 Framing(#[from] FramingError),
40 #[error("failed to decode HandoffOffer Frame: {0}")]
42 DecodeFrame(prost::DecodeError),
43 #[error("failed to decode HandoffOffer payload: {0}")]
45 DecodePayload(prost::DecodeError),
46 #[error("unexpected HandoffOffer frame: {0}")]
48 UnexpectedFrame(&'static str),
49}
50
51impl From<std::io::Error> for BackendHandoffWireError {
52 fn from(error: std::io::Error) -> Self {
53 Self::Framing(FramingError::Io(error))
54 }
55}
56
57pub fn read_handoff_offer<S: Read>(
59 stream: &mut S,
60) -> Result<HandoffOffer, BackendHandoffWireError> {
61 let bytes = read_frame(stream)?;
62 let frame = Frame::decode(bytes.as_slice()).map_err(BackendHandoffWireError::DecodeFrame)?;
63 validate_handoff_frame(&frame, FrameKind::Request)
64 .map_err(BackendHandoffWireError::UnexpectedFrame)?;
65 let offer = HandoffOffer::decode(frame.payload.as_slice())
66 .map_err(BackendHandoffWireError::DecodePayload)?;
67 if frame.request_id != offer.correlation_id {
68 return Err(BackendHandoffWireError::UnexpectedFrame(
69 "frame request_id does not match HandoffOffer correlation_id",
70 ));
71 }
72 Ok(offer)
73}
74
75pub fn read_handoff_offer_with_deadline(
78 stream: &mut interprocess::local_socket::Stream,
79 deadline: Instant,
80) -> Result<HandoffOffer, BackendHandoffWireError> {
81 with_nonblocking_deadline(stream, deadline, |stream| read_handoff_offer(stream))
82}
83
84pub fn respond_to_handoff_offer<S: Write>(
93 stream: &mut S,
94 pending_tokens: &mut HandoffTokenStore,
95 expected_token: HandoffToken,
96 offer: HandoffOffer,
97 now: Instant,
98) -> Result<HandoffAcceptance<HandoffOffer>, BackendHandoffWireError> {
99 let presented_token = offer.token.clone();
100 let correlation_id = offer.correlation_id;
101 let payload = HandedOffPayload::new(expected_token, presented_token.clone(), offer);
102 let acceptance = accept_handed_off(pending_tokens, payload, now);
103
104 let ack = match &acceptance {
105 HandoffAcceptance::Accepted(_) => HandoffAck {
106 token: presented_token,
107 accepted: true,
108 error_detail: String::new(),
109 correlation_id,
110 },
111 HandoffAcceptance::Rejected(rejected) => HandoffAck {
112 token: presented_token,
113 accepted: false,
114 error_detail: rejected.reason.to_string(),
115 correlation_id,
116 },
117 };
118 write_handoff_ack(stream, &ack)?;
119 Ok(acceptance)
120}
121
122pub fn write_handoff_ack<S: Write>(
124 stream: &mut S,
125 ack: &HandoffAck,
126) -> Result<(), BackendHandoffWireError> {
127 let frame = handoff_ack_frame(ack);
128 let mut bytes = Vec::with_capacity(64);
129 frame
130 .encode(&mut bytes)
131 .expect("prost encoding Frame into Vec cannot fail because Vec writes are infallible");
132 write_frame(stream, &bytes)?;
133 Ok(())
134}
135
136#[deprecated(
144 since = "4.6.2",
145 note = "use serve_handoff_offer_with_deadline for production local sockets"
146)]
147pub fn serve_handoff_offer<S: Read + Write>(
148 stream: &mut S,
149 pending_tokens: &mut HandoffTokenStore,
150 expected_token: HandoffToken,
151 now: Instant,
152) -> Result<HandoffAcceptance<HandoffOffer>, BackendHandoffWireError> {
153 serve_handoff_offer_inner(stream, pending_tokens, expected_token, now)
154}
155
156fn serve_handoff_offer_inner<S: Read + Write>(
157 stream: &mut S,
158 pending_tokens: &mut HandoffTokenStore,
159 expected_token: HandoffToken,
160 now: Instant,
161) -> Result<HandoffAcceptance<HandoffOffer>, BackendHandoffWireError> {
162 let offer = read_handoff_offer(stream)?;
163 respond_to_handoff_offer(stream, pending_tokens, expected_token, offer, now)
164}
165
166pub fn serve_handoff_offer_with_deadline(
185 stream: &mut interprocess::local_socket::Stream,
186 pending_tokens: &mut HandoffTokenStore,
187 expected_token: HandoffToken,
188 now: Instant,
189 deadline: Instant,
190) -> Result<HandoffAcceptance<HandoffOffer>, BackendHandoffWireError> {
191 let offer = read_handoff_offer_with_deadline(stream, deadline)?;
192 respond_to_handoff_offer(stream, pending_tokens, expected_token, offer, now)
193}