heddle-thread-api 0.28.9

Native Thread clients and durable peer replication over Iroh
Documentation
use std::{
    future::Future,
    pin::pin,
    sync::{
        Arc,
        atomic::{AtomicUsize, Ordering},
    },
    task::{Context, Poll, Waker},
};

use api::v2::client::{ClientError, MessageReader, MessageWriter, Rpc, RpcTransport};
use prost::Message;

use crate::{Remote, contract::*, rpc, transport::Error};

struct Empty;
#[test]
fn sync_call_context_and_openings_share_the_disabled_cutover_switch() {
    let call_protocol = |method: &api::v2::MethodDescriptor| {
        super::call_protocol(method.path, !method.mandatory_features.is_empty())
    };
    assert_eq!(super::sync_protocol(), None);
    assert_eq!(call_protocol(rpc::SyncServiceFetch::METHOD), None);
    assert_eq!(
        super::call_protocol(rpc::SyncServiceFetch::METHOD.path, true),
        None,
        "API metadata alone cannot activate the coordinated Sync cutover"
    );
    assert_eq!(call_protocol(rpc::SyncServicePublishContent::METHOD), None);
    assert_eq!(call_protocol(rpc::SyncServiceReplicateThread::METHOD), None);
    assert_eq!(
        call_protocol(rpc::IntegrationServicePrepareImportJob::METHOD),
        Some(super::protocol())
    );
    assert_eq!(
        call_protocol(rpc::EndpointServiceDescribeEndpoint::METHOD),
        None
    );
}
#[test]
fn a_complete_public_bundle_has_a_capable_transport_control() {
    let fixture: serde_json::Value =
        serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
            .expect("alpha.21 vectors");
    let bytes = hex::decode(
        fixture["wire_vectors"]["complete_export"]["wire_hex"]
            .as_str()
            .expect("wire bytes"),
    )
    .expect("hex");
    let bundle = ImportPublicProofBundleV1::decode(bytes.as_slice()).expect("complete export");
    let open = ReplicationOpen {
        protocol: Some(super::protocol()),
        import_authority: Some(bundle.clone()),
        ..Default::default()
    };
    super::replication_open(&open).expect("supported complete carrier");
    super::operations(&ReplicationOperations {
        import_authority: Some(bundle),
        ..Default::default()
    })
    .expect("structural complete closure");
    let mut incomplete = open.clone();
    incomplete
        .import_authority
        .as_mut()
        .expect("bundle")
        .genesis_witnesses
        .clear();
    assert!(super::replication_open(&incomplete).is_err());
    let mut old = open.clone();
    old.protocol = None;
    assert!(super::replication_open(&old).is_err());
    super::replication_open(&open).expect("complete capable control remains accepted");
}

#[test]
fn boundary_acceptance_is_carried_with_its_exact_api_binding() {
    let fixture: serde_json::Value =
        serde_json::from_str(include_str!("../../tests/fixtures/hybrid-alpha33.json"))
            .expect("published vectors");
    let record = |section: &str, name: &str| {
        hex::decode(
            fixture[section][name]["wire_hex"]
                .as_str()
                .expect("wire bytes"),
        )
        .expect("hex")
    };
    let mut bundle =
        ImportPublicProofBundleV1::decode(record("wire_vectors", "complete_export").as_slice())
            .expect("complete export");
    let payload = ImportGenesisWitnessV1::decode(
        record("wire_vectors", "boundary_genesis_payload").as_slice(),
    )
    .expect("exact boundary payload");
    let statement = api::heddle::api::common::SignedHostedWitnessStatementV1::decode(
        record("signed_vectors", "boundary_genesis_statement").as_slice(),
    )
    .expect("original boundary statement");
    let original = bundle
        .genesis_witnesses
        .iter_mut()
        .find(|p| p.original_genesis == payload.original_genesis)
        .expect("selected original genesis");
    let old_payload = api::hybrid_codec::canonical(original).expect("original payload");
    *original = payload.clone();
    bundle.statements.retain(|s| {
        !s.body
            .as_ref()
            .is_some_and(|s| s.purpose == 1 && s.canonical_payload == old_payload)
    });
    bundle.statements.push(statement.clone());
    let open = ReplicationOpen {
        protocol: Some(super::protocol()),
        import_authority: Some(bundle.clone()),
        ..Default::default()
    };
    super::replication_open(&open).expect("complete boundary evidence is transportable");
    api::import_authority::verify_witness_payload(
        statement.body.as_ref().expect("statement"),
        api::import_authority::WitnessPayload::Genesis(&payload),
    )
    .expect("API verifies exact originals, manifest, intent and receipts");
    let mut missing = bundle;
    missing
        .genesis_witnesses
        .iter_mut()
        .for_each(|p| p.boundary_acceptance = None);
    assert!(super::bundle(Some(&missing)).is_err());
    super::replication_open(&open).expect("unchanged boundary control");
}
impl MessageReader for Empty {
    type Error = Error;
    async fn next(&mut self) -> Result<Option<Vec<u8>>, Error> {
        Ok(None)
    }
    fn cancel(&mut self) {}
}
impl MessageWriter for Empty {
    type Error = Error;
    async fn send(&mut self, _: Vec<u8>) -> Result<(), Error> {
        Ok(())
    }
    async fn finish(&mut self) -> Result<(), Error> {
        Ok(())
    }
    fn abort(&mut self) {}
}
struct Peer {
    description: DescribeEndpointResponse,
    calls: Arc<AtomicUsize>,
}
impl RpcTransport for Peer {
    type Error = Error;
    type Reader = Empty;
    type Writer = Empty;
    async fn unary(
        &self,
        method: &'static api::v2::MethodDescriptor,
        _: Vec<u8>,
    ) -> Result<Vec<u8>, Error> {
        self.calls.fetch_add(1, Ordering::Relaxed);
        if method.path == rpc::EndpointServiceDescribeEndpoint::METHOD.path {
            return Ok(self.description.encode_to_vec());
        }
        Ok(PrepareImportJobResponse::default().encode_to_vec())
    }
    async fn observe(
        &self,
        _: &'static api::v2::MethodDescriptor,
        _: Vec<u8>,
    ) -> Result<Empty, Error> {
        Ok(Empty)
    }
    async fn exchange(
        &self,
        _: &'static api::v2::MethodDescriptor,
        _: Vec<u8>,
    ) -> Result<(Empty, Empty), Error> {
        Ok((Empty, Empty))
    }
}
fn ready<F: Future>(future: F) -> F::Output {
    match pin!(future)
        .as_mut()
        .poll(&mut Context::from_waker(Waker::noop()))
    {
        Poll::Ready(result) => result,
        Poll::Pending => panic!("immediate fixture"),
    }
}

#[test]
fn discovery_retains_peer_support_and_rejects_old_peers_before_the_gated_call() {
    use api::heddle::api::common::ProtocolCompatibility;
    for (protocol, accepted) in [
        (None, false),
        (Some(ProtocolCompatibility::default()), false),
        (
            Some(ProtocolCompatibility {
                protocol_version: 2,
                mandatory_features: vec![1, 1],
            }),
            false,
        ),
        (
            Some(ProtocolCompatibility {
                protocol_version: 2,
                mandatory_features: vec![1, 2],
            }),
            false,
        ),
        (Some(super::protocol()), true),
    ] {
        let calls = Arc::new(AtomicUsize::new(0));
        let peer = Peer {
            calls: calls.clone(),
            description: DescribeEndpointResponse {
                endpoint: Some(EndpointRef {
                    public_key: vec![7; 32],
                    kind: EndpointKind::Weft as i32,
                }),
                implemented_methods: vec![
                    rpc::IntegrationServicePrepareImportJob::METHOD.path.into(),
                ],
                supported_packages: vec!["heddle.api.v1alpha2".into()],
                protocol,
                ..Default::default()
            },
        };
        let remote = ready(Remote::discover(peer, [7; 32], EndpointKind::Weft))
            .expect("authenticated discovery");
        let result = ready(remote.api.call::<rpc::IntegrationServicePrepareImportJob>(
            &PrepareImportJobRequest {
                client_operation_id: "stable-operation".into(),
                ..Default::default()
            },
        ));
        if accepted {
            result.expect("capable peer");
        } else {
            assert!(matches!(result, Err(ClientError::Protocol(_))));
        }
        assert_eq!(
            calls.load(Ordering::Relaxed),
            if accepted { 2 } else { 1 },
            "gate must precede request transport"
        );
    }
}

#[test]
fn sync_mandatory_gate_stays_off_and_optional_negotiation_is_exact() {
    const { assert!(!super::SYNC_MANDATORY_GATE) };
    for method in [
        rpc::SyncServiceFetch::METHOD,
        rpc::SyncServicePublishContent::METHOD,
        rpc::SyncServiceReplicateThread::METHOD,
    ] {
        assert!(method.mandatory_features.is_empty());
    }
    super::negotiated(None, None).expect("ordinary Sync");
    super::negotiated(Some(&super::protocol()), Some(&super::protocol()))
        .expect("capable Sync path");
    assert!(super::negotiated(Some(&super::protocol()), None).is_err());
}