Skip to main content

running_process/broker/protocol_v2/
mod.rs

1//! v2 broker protocol module.
2//!
3//! Houses the prost-generated types for the `running_process.broker.v2`
4//! package — currently the `ServiceDefinition` envelope and the
5//! `HttpServerCapability` optional sub-message introduced in #483.
6//!
7//! v2 runs in parallel with v1 (`super::protocol`) through the broker
8//! v2 rollout; v1's types are FROZEN FOREVER (#228) so all new
9//! capability fields land here instead.
10
11pub use running_process_protocol::broker::v2::*;
12
13impl running_process_protocol::SessionStartEnvironmentPolicy for crate::EnvironmentPolicy {
14    fn session_start_wire_fields(self) -> (i32, bool) {
15        let policy = match self {
16            crate::EnvironmentPolicy::Auto => crate::EnvironmentPolicy::Inherit,
17            explicit => explicit,
18        };
19        (
20            policy.wire_value().expect("resolved policy"),
21            policy.legacy_clear_fallback().expect("resolved policy"),
22        )
23    }
24}
25
26pub use crate::daemon_registration_v2::{
27    service_definition_dir_v2, service_definition_path_v2, write_service_definition_v2,
28    ServiceDefinitionBuilder, SERVICE_DEF_V2_EXTENSION,
29};
30
31mod manifest_io;
32pub use manifest_io::{
33    central_manifest_path_v2, central_registry_dir_v2, write_to_central_in_dir_v2,
34    write_to_central_v2, write_to_root_v2, CacheManifestBuilder, BROKER_ENVELOPE_VERSION_V2,
35    CENTRAL_MANIFEST_EXTENSION_V2, ROOT_MANIFEST_FILE_V2,
36};
37
38pub mod backend_handle;
39
40pub mod client_compat;
41
42mod loader;
43pub use loader::{ServiceDefinitionLoader, ServiceDefinitionScanEntry};
44
45#[cfg(test)]
46mod tests {
47    use super::*;
48    use prost::Message;
49
50    /// `ServiceDefinition` round-trips with no HTTP capability — the
51    /// optional field is absent on both sides.
52    #[test]
53    fn service_definition_without_http_round_trips() {
54        let original = ServiceDefinition {
55            service_name: "zccache".to_owned(),
56            http_server: None,
57            ..Default::default()
58        };
59
60        let bytes = original.encode_to_vec();
61        let decoded =
62            ServiceDefinition::decode(bytes.as_slice()).expect("encoded ServiceDefinition decodes");
63
64        assert_eq!(decoded.service_name, "zccache");
65        assert!(decoded.http_server.is_none());
66    }
67
68    /// soldr#2365 Phase 3: every `SessionFrame` variant survives a prost
69    /// round-trip byte-exactly, including raw non-UTF-8 stdio payloads (the
70    /// proxy data plane must be byte-transparent, not text-normalizing).
71    #[test]
72    fn session_frame_each_variant_round_trips() {
73        let raw: Vec<u8> = vec![0x00, 0xff, 0x80, b'\n', 0x1b];
74        let cases = [
75            session_frame::Kind::Stdin(raw.clone()),
76            session_frame::Kind::StdinEof(true),
77            session_frame::Kind::Stdout(raw.clone()),
78            session_frame::Kind::Stderr(raw.clone()),
79            session_frame::Kind::Exit(SessionExit {
80                code: 101,
81                signal: 0,
82                metadata: Default::default(),
83            }),
84            session_frame::Kind::Start(SessionStart {
85                program: "rustc".to_owned(),
86                args: vec!["--edition".to_owned(), "2021".to_owned()],
87                cwd: "/work".to_owned(),
88                env: vec![SessionEnvVar {
89                    key: "CARGO".to_owned(),
90                    value: "1".to_owned(),
91                }],
92                clear_inherited_env: true,
93                environment_policy: 3,
94            }),
95        ];
96        for kind in cases {
97            let original = SessionFrame {
98                kind: Some(kind.clone()),
99            };
100            let decoded = SessionFrame::decode(original.encode_to_vec().as_slice())
101                .expect("SessionFrame decodes");
102            assert_eq!(decoded, original, "variant did not round-trip: {kind:?}");
103        }
104    }
105
106    #[test]
107    fn session_start_captures_context_and_dual_writes_policy() {
108        let start = SessionStart::from_current_process("rustc", ["--version"], "work");
109        assert_eq!(start.environment_policy, 3);
110        assert!(start.clear_inherited_env);
111        assert!(!start.env.is_empty());
112
113        let inherit = start.with_environment_policy(crate::EnvironmentPolicy::Inherit);
114        assert_eq!(inherit.environment_policy, 1);
115        assert!(!inherit.clear_inherited_env);
116    }
117
118    /// Signal death carries a non-zero `signal`; consumers read it first.
119    #[test]
120    fn session_exit_signal_death_round_trips() {
121        let original = SessionExit {
122            code: -1,
123            signal: 9,
124            // Non-empty so the opaque metadata map (soldr#2365 Q3) is exercised
125            // across the wire, not just defaulted away.
126            metadata: [
127                ("cache_outcome".to_owned(), "miss".to_owned()),
128                ("compile_id".to_owned(), "compile-934".to_owned()),
129            ]
130            .into_iter()
131            .collect(),
132        };
133        let decoded =
134            SessionExit::decode(original.encode_to_vec().as_slice()).expect("SessionExit decodes");
135        assert_eq!(decoded, original);
136        assert_ne!(decoded.signal, 0, "signal death must be distinguishable");
137    }
138
139    /// The session lane's payload-protocol id is distinct from every other
140    /// registered broker protocol (the registry's pairwise-distinct invariant).
141    #[test]
142    fn session_payload_protocol_is_distinct() {
143        use crate::broker::protocol::{
144            ADMIN_PAYLOAD_PROTOCOL, BACKEND_HANDLE_PROBE_PAYLOAD_PROTOCOL,
145            CONTROL_PAYLOAD_PROTOCOL, HANDOFF_PAYLOAD_PROTOCOL, SESSION_PAYLOAD_PROTOCOL,
146        };
147        assert_eq!(SESSION_PAYLOAD_PROTOCOL, 0x5350);
148        for other in [
149            CONTROL_PAYLOAD_PROTOCOL,
150            ADMIN_PAYLOAD_PROTOCOL,
151            BACKEND_HANDLE_PROBE_PAYLOAD_PROTOCOL,
152            HANDOFF_PAYLOAD_PROTOCOL,
153        ] {
154            assert_ne!(SESSION_PAYLOAD_PROTOCOL, other);
155        }
156    }
157
158    /// `ServiceDefinition` round-trips with an `HttpServerCapability`
159    /// populated — all three fields survive.
160    #[test]
161    fn service_definition_with_http_round_trips() {
162        let original = ServiceDefinition {
163            service_name: "fbuild".to_owned(),
164            http_server: Some(HttpServerCapability {
165                bind_addr: "127.0.0.1".to_owned(),
166                health_path: "/healthz".to_owned(),
167                display_name: "fbuild status".to_owned(),
168            }),
169            ..Default::default()
170        };
171
172        let bytes = original.encode_to_vec();
173        let decoded =
174            ServiceDefinition::decode(bytes.as_slice()).expect("encoded ServiceDefinition decodes");
175
176        let cap = decoded
177            .http_server
178            .expect("http_server survives round-trip");
179        assert_eq!(decoded.service_name, "fbuild");
180        assert_eq!(cap.bind_addr, "127.0.0.1");
181        assert_eq!(cap.health_path, "/healthz");
182        assert_eq!(cap.display_name, "fbuild status");
183    }
184
185    /// Empty `HttpServerCapability` survives a round-trip — defaults are
186    /// applied by the loader/consumer, not by the proto encoder.
187    #[test]
188    fn http_server_capability_empty_defaults_survive_round_trip() {
189        let original = ServiceDefinition {
190            service_name: "minimal".to_owned(),
191            http_server: Some(HttpServerCapability::default()),
192            ..Default::default()
193        };
194
195        let bytes = original.encode_to_vec();
196        let decoded =
197            ServiceDefinition::decode(bytes.as_slice()).expect("encoded ServiceDefinition decodes");
198
199        let cap = decoded
200            .http_server
201            .expect("http_server survives round-trip");
202        assert!(cap.bind_addr.is_empty());
203        assert!(cap.health_path.is_empty());
204        assert!(cap.display_name.is_empty());
205    }
206
207    /// Slice 22 (zackees/zccache#782): the launcher / isolation fields
208    /// ported from v1 round-trip cleanly. Pins every new field plus the
209    /// `BrokerIsolation` enum mapping so a future proto regression
210    /// surfaces here instead of at the first downstream loader.
211    #[test]
212    fn service_definition_v1_fields_round_trip() {
213        use std::collections::HashMap;
214        let mut labels = HashMap::new();
215        labels.insert("env".to_owned(), "prod".to_owned());
216        labels.insert("deploy".to_owned(), "blue".to_owned());
217
218        let original = ServiceDefinition {
219            service_name: "zccache".to_owned(),
220            binary_path: "/usr/local/bin/zccache-daemon".to_owned(),
221            isolation: BrokerIsolation::SharedBroker as i32,
222            explicit_instance: String::new(),
223            per_version_binary_dir: "/usr/local/bin".to_owned(),
224            min_version: "1.0.0".to_owned(),
225            version_allow_list: vec!["1.12.9".to_owned(), "1.13.0".to_owned()],
226            labels,
227            http_server: None,
228        };
229
230        let bytes = original.encode_to_vec();
231        let decoded =
232            ServiceDefinition::decode(bytes.as_slice()).expect("encoded ServiceDefinition decodes");
233
234        assert_eq!(decoded.service_name, "zccache");
235        assert_eq!(decoded.binary_path, "/usr/local/bin/zccache-daemon");
236        assert_eq!(decoded.isolation, BrokerIsolation::SharedBroker as i32);
237        assert!(decoded.explicit_instance.is_empty());
238        assert_eq!(decoded.per_version_binary_dir, "/usr/local/bin");
239        assert_eq!(decoded.min_version, "1.0.0");
240        assert_eq!(
241            decoded.version_allow_list,
242            vec!["1.12.9".to_owned(), "1.13.0".to_owned()]
243        );
244        assert_eq!(decoded.labels.len(), 2);
245        assert_eq!(decoded.labels.get("env"), Some(&"prod".to_owned()));
246        assert_eq!(decoded.labels.get("deploy"), Some(&"blue".to_owned()));
247        assert!(decoded.http_server.is_none());
248    }
249
250    /// Slice 22 (zackees/zccache#782): every `BrokerIsolation` enum
251    /// variant survives the round-trip. Pins the proto-int mapping so
252    /// future variant additions get an explicit failure here instead
253    /// of misclassifying as `PrivateBroker` (the proto3 zero value).
254    #[test]
255    fn broker_isolation_enum_values_round_trip() {
256        for iso in [
257            BrokerIsolation::PrivateBroker,
258            BrokerIsolation::SharedBroker,
259            BrokerIsolation::ExplicitInstance,
260        ] {
261            let original = ServiceDefinition {
262                service_name: format!("svc-{}", iso as i32),
263                isolation: iso as i32,
264                ..Default::default()
265            };
266            let bytes = original.encode_to_vec();
267            let decoded = ServiceDefinition::decode(bytes.as_slice())
268                .expect("encoded ServiceDefinition decodes");
269            assert_eq!(decoded.isolation, iso as i32, "round-trip of {iso:?}");
270        }
271    }
272
273    /// Slice 22: `explicit_instance` is only meaningful when
274    /// `isolation == ExplicitInstance`, but the proto encoder doesn't
275    /// enforce the gating — the consumer (broker / loader) does. Pin
276    /// that the field round-trips regardless of the isolation value
277    /// so a future broker policy change can rely on the bytes round-tripping
278    /// faithfully even for "invalid" combinations.
279    #[test]
280    fn service_definition_explicit_instance_round_trips_with_any_isolation() {
281        for iso in [
282            BrokerIsolation::PrivateBroker,
283            BrokerIsolation::SharedBroker,
284            BrokerIsolation::ExplicitInstance,
285        ] {
286            let original = ServiceDefinition {
287                service_name: "svc".to_owned(),
288                isolation: iso as i32,
289                explicit_instance: "ci-trusted".to_owned(),
290                ..Default::default()
291            };
292            let bytes = original.encode_to_vec();
293            let decoded = ServiceDefinition::decode(bytes.as_slice())
294                .expect("encoded ServiceDefinition decodes");
295            assert_eq!(
296                decoded.explicit_instance, "ci-trusted",
297                "explicit_instance must survive round-trip even with isolation={iso:?}"
298            );
299        }
300    }
301
302    /// `BackendHttpReady` carries the daemon's OS-allocated port back to
303    /// the broker; encodes/decodes without loss.
304    #[test]
305    fn backend_http_ready_round_trips() {
306        let original = BackendHttpReady { port: 49_152 };
307
308        let bytes = original.encode_to_vec();
309        let decoded = BackendHttpReady::decode(bytes.as_slice()).expect("BackendHttpReady decodes");
310
311        assert_eq!(decoded.port, 49_152);
312    }
313
314    /// `GetBrokerHttpEndpointRequest` is an empty marker; encoding +
315    /// decoding it produces the same default-constructed message.
316    #[test]
317    fn get_broker_http_endpoint_request_round_trips_empty() {
318        let original = GetBrokerHttpEndpointRequest::default();
319
320        let bytes = original.encode_to_vec();
321        let decoded = GetBrokerHttpEndpointRequest::decode(bytes.as_slice())
322            .expect("GetBrokerHttpEndpointRequest decodes");
323
324        assert_eq!(decoded, GetBrokerHttpEndpointRequest::default());
325    }
326
327    /// `GetBrokerHttpEndpointResponse` round-trips both fields (port + pid).
328    #[test]
329    fn get_broker_http_endpoint_response_round_trips() {
330        let original = GetBrokerHttpEndpointResponse {
331            port: 8765,
332            pid: 12_345,
333        };
334
335        let bytes = original.encode_to_vec();
336        let decoded = GetBrokerHttpEndpointResponse::decode(bytes.as_slice())
337            .expect("GetBrokerHttpEndpointResponse decodes");
338
339        assert_eq!(decoded.port, 8765);
340        assert_eq!(decoded.pid, 12_345);
341    }
342}