heddle_thread_api/
creation.rs1use api::v2::client::{ClientError, RpcTransport};
4use crypto::Signer;
5use heddle_object_model::object::{OperationId, thread_replication::ThreadGenesis};
6
7use crate::{Remote, contract::*, replication::opening, rpc, transport::Error};
8
9pub struct ThreadCreation {
10 request: StartThreadRequest,
11 reference: ThreadRef,
12}
13
14impl ThreadCreation {
15 pub fn sign(
16 operation_id: impl Into<String>,
17 genesis: &ThreadGenesis,
18 signer: &impl Signer,
19 ) -> Result<Self, Error> {
20 let operation_id = operation_id.into();
21 validate_operation(&operation_id)?;
22 Self::from_signed(operation_id, opening::sign_genesis(genesis, signer)?)
23 }
24
25 pub fn sign_with_authority(
26 operation_id: impl Into<String>,
27 genesis: &ThreadGenesis,
28 signer: &impl Signer,
29 creator_authority: Vec<u8>,
30 ) -> Result<Self, Error> {
31 Self::from_signed_with_authority(
32 operation_id,
33 opening::sign_genesis(genesis, signer)?,
34 creator_authority,
35 )
36 }
37
38 pub fn from_signed(
40 operation_id: impl Into<String>,
41 record: SignedRecord,
42 ) -> Result<Self, Error> {
43 Self::from_signed_with_authority(operation_id, record, Vec::new())
44 }
45
46 pub fn from_signed_with_authority(
49 operation_id: impl Into<String>,
50 record: SignedRecord,
51 creator_authority: Vec<u8>,
52 ) -> Result<Self, Error> {
53 let operation_id = operation_id.into();
54 validate_operation(&operation_id)?;
55 let genesis = ThreadGenesis::decode(&record.canonical_record)
56 .map_err(|_| Error::Protocol("invalid canonical Thread genesis"))?;
57 use heddle_object_model::object::thread_replication::GenesisOwner;
58 match &genesis.owner {
59 GenesisOwner::LocalKey(_) if !creator_authority.is_empty() => {
60 return Err(Error::Protocol(
61 "local-key ownership requires an explicit claim, not an account proof on upload",
62 ));
63 }
64 GenesisOwner::Account(_) if creator_authority.is_empty() => {
65 return Err(Error::Protocol(
66 "account-owned genesis requires original creator authority",
67 ));
68 }
69 _ => {}
70 }
71 if creator_authority.len() > 64 * 1024 {
72 return Err(Error::Protocol("creator authority exceeds bound"));
73 }
74 let reference = ThreadRef {
75 spool: Some(SpoolRef {
76 id: genesis.spool.clone(),
77 }),
78 id: Some(ThreadId {
79 value: genesis
80 .id()
81 .map_err(|_| Error::Protocol("invalid Thread identity"))?
82 .as_bytes()
83 .to_vec(),
84 }),
85 };
86 opening::verify_genesis(&record, &reference)?;
87 Ok(Self {
88 request: StartThreadRequest {
89 native_genesis_authority: None,
90 client_operation_id: operation_id,
91 spool: reference.spool.clone(),
92 thread_genesis: Some(record),
93 creator_authority,
94 },
95 reference,
96 })
97 }
98
99 pub fn with_native_authority(
101 mut self,
102 binding: SignedNativeGenesisAuthorityV1,
103 ) -> Result<Self, Error> {
104 api::native_witness::verify_genesis_authority(
105 &binding,
106 self.request
107 .thread_genesis
108 .as_ref()
109 .ok_or(Error::Protocol("genesis absent"))?,
110 &self.request.creator_authority,
111 )
112 .map_err(|_| {
113 Error::Protocol("native creator binding does not match genesis and authority")
114 })?;
115 self.request.native_genesis_authority = Some(binding);
116 Ok(self)
117 }
118
119 pub fn request(&self) -> &StartThreadRequest {
120 &self.request
121 }
122 pub fn reference(&self) -> &ThreadRef {
123 &self.reference
124 }
125 pub fn genesis_record(&self) -> ThreadGenesisRecord {
126 ThreadGenesisRecord {
127 native_genesis_authority: self.request.native_genesis_authority.clone(),
128 boundary_acceptances: Vec::new(),
129 ownership_claims: vec![],
130 ownership_claim_admissions: vec![],
131 ownership_resolutions: vec![],
132 ownership_resolution_admissions: vec![],
133 genesis: self.request.thread_genesis.clone(),
134 creator_authority: self.request.creator_authority.clone(),
135 admission: None,
136 }
137 }
138}
139
140fn validate_operation(operation_id: &str) -> Result<(), Error> {
141 operation_id
142 .parse::<OperationId>()
143 .map(|_| ())
144 .map_err(|_| Error::Protocol("client operation ID must be a UUID"))
145}
146
147impl<T: RpcTransport<Error = Error>> Remote<T> {
148 pub async fn start_thread(
151 &self,
152 creation: &ThreadCreation,
153 ) -> Result<ThreadMutationResponse, ClientError<Error>> {
154 if creation.request().native_genesis_authority.is_some() {
155 api::import_authority::require_hybrid_peer(self.description.protocol.as_ref())
156 .map_err(|_| {
157 ClientError::Transport(Error::Protocol("native genesis requires a HYBRID peer"))
158 })?;
159 }
160 if !self
161 .description
162 .understood_signed_record_formats
163 .iter()
164 .any(|format| format == heddle_object_model::object::thread_replication::GENESIS_FORMAT)
165 {
166 return Err(ClientError::Transport(Error::Protocol(
167 "endpoint does not understand Thread genesis",
168 )));
169 }
170 let response = self
171 .api
172 .call::<rpc::ThreadServiceStartThread>(creation.request())
173 .await?;
174 let receipt = response
175 .receipt
176 .as_ref()
177 .ok_or(ClientError::Transport(Error::Protocol(
178 "Thread creation response has no receipt",
179 )))?;
180 if receipt.client_operation_id != creation.request.client_operation_id
181 || receipt.endpoint != self.description.endpoint
182 || receipt.outcome.is_none()
183 || response
184 .thread
185 .as_ref()
186 .is_some_and(|thread| thread.r#ref.as_ref() != Some(creation.reference()))
187 || (matches!(receipt.outcome, Some(mutation_receipt::Outcome::Applied(_)))
188 && response.thread.is_none())
189 {
190 return Err(ClientError::Transport(Error::Protocol(
191 "Thread creation response does not match the command",
192 )));
193 }
194 Ok(response)
195 }
196}