1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
//! `GetBrokerHttpEndpoint` RPC dispatch (slice 6 of #488).
//!
//! The CLI calls `GetBrokerHttpEndpoint` over the v2 broker control
//! channel to discover the broker's HTTP endpoint (per #483 §4 — the
//! single discovery surface). This module implements the broker side
//! of that RPC: given the broker's currently-resolved HTTP port + its
//! own pid, build a `GetBrokerHttpEndpointResponse` and serialize it.
//!
//! The real plumbing (read incoming frame → dispatch on payload type →
//! write response frame) lives in the broker's connection loop, which
//! is filled in by later slices. This slice exposes the typed
//! request/response handler so subsequent slices have a pinned API to
//! call.
use prost::Message;
use crate::broker::protocol_v2::{GetBrokerHttpEndpointRequest, GetBrokerHttpEndpointResponse};
/// In-broker resolved HTTP endpoint state (set at boot per #483 §3 via
/// `BrokerHttpPort::resolve(config, env)`).
#[derive(Debug, Clone, Copy)]
pub struct BrokerHttpEndpoint {
/// The port the broker's own HTTP server bound. Slice 7 actually
/// binds it; before then the broker can stub this to its
/// configured-static port for early consumer testing.
pub port: u16,
/// The broker's process id. Used by consumers to disambiguate a
/// fresh response from a stale one mid-restart (#483 §4 rationale).
pub pid: u32,
}
impl BrokerHttpEndpoint {
/// Build a `GetBrokerHttpEndpointResponse` carrying this endpoint.
pub fn to_response(self) -> GetBrokerHttpEndpointResponse {
GetBrokerHttpEndpointResponse {
port: self.port as u32,
pid: self.pid,
}
}
}
/// Errors from [`decode_request_and_dispatch`].
#[derive(Debug, thiserror::Error)]
pub enum GetHttpEndpointError {
/// The incoming frame body did not decode as `GetBrokerHttpEndpointRequest`.
#[error("decode GetBrokerHttpEndpointRequest: {0}")]
Decode(#[from] prost::DecodeError),
/// Encoding the response failed.
#[error("encode GetBrokerHttpEndpointResponse: {0}")]
Encode(#[from] prost::EncodeError),
}
/// Decode an incoming `GetBrokerHttpEndpointRequest` frame body and
/// produce a serialized `GetBrokerHttpEndpointResponse` body the
/// connection loop can write back via `protocol::write_frame`.
///
/// The request currently has no fields (`GetBrokerHttpEndpointRequest`
/// is an empty marker per #483 §4) — decoding is purely validation
/// that the peer sent a structurally well-formed proto message of the
/// expected type.
pub fn decode_request_and_dispatch(
request_body: &[u8],
endpoint: BrokerHttpEndpoint,
) -> Result<Vec<u8>, GetHttpEndpointError> {
let _request = GetBrokerHttpEndpointRequest::decode(request_body)?;
let response = endpoint.to_response();
let mut body = Vec::with_capacity(response.encoded_len());
response.encode(&mut body)?;
Ok(body)
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn to_response_carries_port_and_pid() {
let resp = BrokerHttpEndpoint {
port: 8765,
pid: 12_345,
}
.to_response();
assert_eq!(resp.port, 8765);
assert_eq!(resp.pid, 12_345);
}
#[test]
fn dispatch_round_trip_with_empty_request() {
let req = GetBrokerHttpEndpointRequest::default();
let mut body = Vec::with_capacity(req.encoded_len());
req.encode(&mut body).expect("encode request");
let resp_body = decode_request_and_dispatch(
&body,
BrokerHttpEndpoint {
port: 4242,
pid: 99_999,
},
)
.expect("dispatch succeeds");
let resp =
GetBrokerHttpEndpointResponse::decode(resp_body.as_slice()).expect("decode response");
assert_eq!(resp.port, 4242);
assert_eq!(resp.pid, 99_999);
}
#[test]
fn dispatch_rejects_malformed_request_body() {
let err = decode_request_and_dispatch(
&[0xFF; 4],
BrokerHttpEndpoint {
port: 4242,
pid: 99_999,
},
)
.expect_err("malformed request body should be rejected");
match err {
GetHttpEndpointError::Decode(_) => {}
other => panic!("expected Decode error, got: {other:?}"),
}
}
}