Skip to main content

heddle_thread_api/
hybrid.rs

1//! HYBRID transport checks. Structural proof closure is preparation only;
2//! authority is resolved independently at each serialized durable mutation.
3use api::heddle::api::common::ProtocolCompatibility;
4
5use crate::contract::{
6    FetchOpen, ImportPublicProofBundleV1, NativePublicProofBundleV1, PublicationReceipt,
7    PublishContentOpen, ReplicationOpen, ReplicationOperations, ReplicationReady, StreamOpen,
8    TransferReady,
9};
10
11#[cfg(feature = "native")]
12pub mod authority;
13#[cfg(all(test, feature = "native"))]
14mod foreign_tests;
15pub mod history;
16#[cfg(test)]
17mod history_tests;
18#[cfg(all(test, feature = "native"))]
19mod native_tests;
20#[cfg(test)]
21mod protocol_tests;
22pub type Rejection = &'static str;
23
24/// Integration routes support HYBRID. Keep this single switch OFF until the
25/// coordinated Sync mandatory cutover in api#307; optional capable paths still
26/// require exact protocol negotiation before accepting any authority bundle.
27pub const SYNC_MANDATORY_GATE: bool = false;
28
29/// All ordinary Sync constructors use the same coordinated cutover switch.
30pub fn sync_protocol() -> Option<ProtocolCompatibility> {
31    SYNC_MANDATORY_GATE.then(protocol)
32}
33
34/// Keep RPC preludes and stream openings on the same cutover switch. Until
35/// api#307, Sync declares no protocol support in its CallContext.
36pub fn call_protocol(method_path: &str, requires_hybrid: bool) -> Option<ProtocolCompatibility> {
37    if method_path
38        .trim_start_matches('/')
39        .starts_with("heddle.api.v1alpha2.SyncService/")
40    {
41        return sync_protocol();
42    }
43    requires_hybrid.then(protocol)
44}
45
46/// A creator binding opts StartThread into HYBRID before request PoP and I/O.
47pub fn native_start_thread(method_path: &str, encoded: &[u8]) -> Result<bool, Rejection> {
48    if method_path.trim_start_matches('/') != "heddle.api.v1alpha2.ThreadService/StartThread" {
49        return Ok(false);
50    }
51    use prost::Message;
52    let request = crate::contract::StartThreadRequest::decode(encoded)
53        .map_err(|_| "invalid StartThread request")?;
54    let Some(binding) = request.native_genesis_authority.as_ref() else {
55        return Ok(false);
56    };
57    api::native_witness::verify_genesis_authority(
58        binding,
59        request
60            .thread_genesis
61            .as_ref()
62            .ok_or("native genesis absent")?,
63        &request.creator_authority,
64    )
65    .map_err(|_| "invalid native creator binding")?;
66    Ok(true)
67}
68
69pub fn protocol() -> ProtocolCompatibility {
70    ProtocolCompatibility {
71        protocol_version: 2,
72        mandatory_features: vec![1],
73    }
74}
75fn check_protocol(value: Option<&ProtocolCompatibility>) -> Result<(), Rejection> {
76    if SYNC_MANDATORY_GATE || value.is_some() {
77        api::import_authority::require_hybrid_peer(value)
78            .map_err(|_| "incompatible HYBRID protocol (api#307 cutover)")?;
79    }
80    Ok(())
81}
82/// Complete structural closure, never signature trust, enrollment or admission.
83pub fn bundle(value: Option<&ImportPublicProofBundleV1>) -> Result<(), Rejection> {
84    if let Some(value) = value {
85        api::import_authority::validate_public_bundle(value)
86            .map_err(|_| "incomplete HYBRID import authority (api#307 cutover)")?;
87    }
88    Ok(())
89}
90pub fn bundles(
91    imported: Option<&ImportPublicProofBundleV1>,
92    native: Option<&NativePublicProofBundleV1>,
93) -> Result<(), Rejection> {
94    api::native_witness::validate_carriers(imported, native).map_err(|_| {
95        if native.is_some() {
96            "incomplete or conflicting HYBRID native authority (api#307 cutover)"
97        } else {
98            "incomplete HYBRID import authority (api#307 cutover)"
99        }
100    })
101}
102fn carrier(
103    protocol: Option<&ProtocolCompatibility>,
104    value: Option<&ImportPublicProofBundleV1>,
105    native: Option<&NativePublicProofBundleV1>,
106) -> Result<(), Rejection> {
107    check_protocol(protocol)?;
108    if value.is_some() || native.is_some() {
109        api::import_authority::require_hybrid_peer(protocol).map_err(
110            |_| "import authority requires negotiated HYBRID protocol (api#307 cutover)",
111        )?;
112    }
113    bundles(value, native)
114}
115/// Bind support to both ends of this exact stream, including its first Ready.
116pub fn negotiated(
117    open: Option<&ProtocolCompatibility>,
118    ready: Option<&ProtocolCompatibility>,
119) -> Result<(), Rejection> {
120    check_protocol(open)?;
121    check_protocol(ready)?;
122    if open != ready {
123        return Err("HYBRID protocol differs from stream opening");
124    }
125    Ok(())
126}
127pub fn operations(batch: &ReplicationOperations) -> Result<(), Rejection> {
128    bundles(
129        batch.import_authority.as_ref(),
130        batch.native_authority.as_ref(),
131    )
132}
133pub fn replication_open(open: &ReplicationOpen) -> Result<(), Rejection> {
134    carrier(
135        open.protocol.as_ref(),
136        open.import_authority.as_ref(),
137        open.native_authority.as_ref(),
138    )
139}
140pub fn replication_ready(ready: &ReplicationReady) -> Result<(), Rejection> {
141    carrier(
142        ready.protocol.as_ref(),
143        ready.import_authority.as_ref(),
144        ready.native_authority.as_ref(),
145    )
146}
147pub fn fetch_open(open: &FetchOpen) -> Result<(), Rejection> {
148    check_protocol(open.protocol.as_ref())
149}
150pub fn transfer_ready(ready: &TransferReady) -> Result<(), Rejection> {
151    carrier(
152        ready.protocol.as_ref(),
153        ready.import_authority.as_ref(),
154        ready.native_authority.as_ref(),
155    )
156}
157pub fn publish_open(open: &PublishContentOpen) -> Result<(), Rejection> {
158    carrier(
159        open.protocol.as_ref(),
160        open.import_authority.as_ref(),
161        open.native_authority.as_ref(),
162    )
163}
164pub fn publication_receipt(receipt: &PublicationReceipt) -> Result<(), Rejection> {
165    bundles(
166        receipt.import_authority.as_ref(),
167        receipt.native_authority.as_ref(),
168    )
169}
170pub fn stream_open(open: &StreamOpen) -> Result<(), Rejection> {
171    check_protocol(open.protocol.as_ref())?;
172    if let Some(set) = &open.witness_set {
173        use prost::Message;
174        api::import_authority::require_hybrid_peer(open.protocol.as_ref())
175            .map_err(|_| "witness set requires HYBRID protocol (api#307 cutover)")?;
176        if set.encoded_len() > api::witness_trust::MAX_SET_BYTES || set.body.is_none() {
177            return Err("invalid HYBRID witness set");
178        }
179    }
180    Ok(())
181}
182
183#[cfg(test)]
184mod tests {
185    use api::heddle::api::common::{ProtocolCompatibility, SignedHostedWitnessSetV1};
186
187    use super::*;
188    use crate::contract::ImportPublicProofBundleV1;
189
190    fn rejected(result: Result<(), Rejection>, field: &str) {
191        let message = result.expect_err(field);
192        assert!(message.contains(field), "{message}");
193        assert!(message.contains("api#307"), "{message}");
194    }
195
196    // Presence alone rejects: an empty-but-present bundle or protocol is still
197    // a claim this peer cannot honour, so the defaults below are deliberate.
198    fn bundle() -> Option<ImportPublicProofBundleV1> {
199        Some(ImportPublicProofBundleV1::default())
200    }
201    fn protocol() -> Option<ProtocolCompatibility> {
202        Some(ProtocolCompatibility::default())
203    }
204
205    #[test]
206    fn absent_hybrid_fields_are_accepted() {
207        operations(&ReplicationOperations::default()).expect("operations");
208        replication_open(&ReplicationOpen::default()).expect("replication open");
209        replication_ready(&ReplicationReady::default()).expect("replication ready");
210        fetch_open(&FetchOpen::default()).expect("fetch open");
211        transfer_ready(&TransferReady::default()).expect("transfer ready");
212        publish_open(&PublishContentOpen::default()).expect("publish open");
213        publication_receipt(&PublicationReceipt::default()).expect("receipt");
214        stream_open(&StreamOpen::default()).expect("stream open");
215    }
216
217    #[test]
218    fn incomplete_import_authority_is_rejected_never_ignored() {
219        rejected(
220            operations(&ReplicationOperations {
221                import_authority: bundle(),
222                ..Default::default()
223            }),
224            "import authority",
225        );
226        rejected(
227            replication_open(&ReplicationOpen {
228                import_authority: bundle(),
229                ..Default::default()
230            }),
231            "import authority",
232        );
233        rejected(
234            replication_ready(&ReplicationReady {
235                import_authority: bundle(),
236                ..Default::default()
237            }),
238            "import authority",
239        );
240        rejected(
241            transfer_ready(&TransferReady {
242                import_authority: bundle(),
243                ..Default::default()
244            }),
245            "import authority",
246        );
247        rejected(
248            publish_open(&PublishContentOpen {
249                import_authority: bundle(),
250                ..Default::default()
251            }),
252            "import authority",
253        );
254        rejected(
255            publication_receipt(&PublicationReceipt {
256                import_authority: bundle(),
257                ..Default::default()
258            }),
259            "import authority",
260        );
261    }
262
263    #[test]
264    fn unsupported_protocol_and_incomplete_witness_sets_are_rejected() {
265        rejected(
266            replication_open(&ReplicationOpen {
267                protocol: protocol(),
268                ..Default::default()
269            }),
270            "protocol",
271        );
272        rejected(
273            replication_ready(&ReplicationReady {
274                protocol: protocol(),
275                ..Default::default()
276            }),
277            "protocol",
278        );
279        rejected(
280            fetch_open(&FetchOpen {
281                protocol: protocol(),
282                ..Default::default()
283            }),
284            "protocol",
285        );
286        rejected(
287            transfer_ready(&TransferReady {
288                protocol: protocol(),
289                ..Default::default()
290            }),
291            "protocol",
292        );
293        rejected(
294            publish_open(&PublishContentOpen {
295                protocol: protocol(),
296                ..Default::default()
297            }),
298            "protocol",
299        );
300        rejected(
301            stream_open(&StreamOpen {
302                protocol: protocol(),
303                ..Default::default()
304            }),
305            "protocol",
306        );
307        rejected(
308            stream_open(&StreamOpen {
309                witness_set: Some(SignedHostedWitnessSetV1::default()),
310                ..Default::default()
311            }),
312            "witness set",
313        );
314    }
315}