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