1use api::v2::client::{ClientError, RpcTransport};
5use prost::Message;
6use tokio::io::{AsyncRead, AsyncReadExt};
7
8use crate::{Remote, contract::*, rpc, transport};
9
10mod batching;
11
12#[cfg(feature = "source-transfer")]
13mod acceptance;
14#[cfg(feature = "source-transfer")]
15pub use acceptance::{
16 PreparedPublication, ProposedAcceptance, PublicationAcceptancePlan, proposed_publication,
17 publication_intent,
18};
19#[cfg(feature = "source-transfer")]
20mod source;
21#[cfg(feature = "source-transfer")]
22mod staging;
23#[cfg(feature = "source-transfer")]
24pub use source::{PublicationOptions, SourceBudget, SourcePack, VisibleSourcePack};
25#[cfg(feature = "source-transfer")]
26pub use staging::{
27 ProposedSourceArtifacts, validate_proposed_source_artifacts, validate_source_artifacts,
28 validate_source_artifacts_with_import_carriers,
29};
30
31#[derive(Clone)]
35pub struct PublicationOriginals {
36 pub geneses: Vec<ThreadGenesisRecord>,
37 pub operations: Vec<ReplicationOperations>,
38}
39impl PublicationOriginals {
40 fn validate_bounds(&self) -> Result<(), Error> {
41 if self.geneses.is_empty()
42 || self.geneses.len() > 128
43 || self.operations.is_empty()
44 || self.operations.len() > 10_000
45 || self
46 .operations
47 .iter()
48 .map(|batch| batch.operations.len())
49 .sum::<usize>()
50 > 10_000
51 || self.operations.iter().any(|batch| {
52 batch.operations.is_empty()
53 || batch.authority_admissions.len() > batch.operations.len()
54 })
55 {
56 return Err(Error::Invalid(
57 "bounded original genesis and source operations required",
58 ));
59 }
60 let mut acceptances = std::collections::BTreeSet::new();
61 for record in self
62 .geneses
63 .iter()
64 .flat_map(|value| &value.boundary_acceptances)
65 .chain(
66 self.operations
67 .iter()
68 .flat_map(|value| &value.boundary_acceptances),
69 )
70 {
71 if record.canonical_record.len() > 96 * 1024
72 || record.signatures.len() != 1
73 || record.signatures[0].signature.len() != 64
74 || record.signatures[0].public_key.len() != 32
75 {
76 return Err(Error::Invalid("boundary evidence shape exceeds bounds"));
77 }
78 acceptances.insert(record.canonical_record.as_slice());
79 if acceptances.len() > 128 {
80 return Err(Error::Invalid("boundary acceptance count exceeded"));
81 }
82 }
83 let mut bytes = 0usize;
84 for length in self
85 .geneses
86 .iter()
87 .map(Message::encoded_len)
88 .chain(self.operations.iter().map(Message::encoded_len))
89 {
90 bytes = bytes
91 .checked_add(length)
92 .ok_or(Error::Invalid("publication original size overflow"))?;
93 if length > 256 * 1024 || bytes > 16 * 1024 * 1024 {
94 return Err(Error::Invalid(
95 "publication original metadata budget exceeded",
96 ));
97 }
98 }
99 if self.geneses.iter().any(|genesis| {
100 genesis.genesis.as_ref().is_none_or(|record| {
101 record.canonical_record.is_empty() || record.signatures.is_empty()
102 })
103 }) || self
104 .operations
105 .iter()
106 .flat_map(|batch| {
107 batch
108 .operations
109 .iter()
110 .chain(&batch.authority_admissions)
111 .chain(&batch.boundary_acceptances)
112 })
113 .any(|record| record.canonical_record.is_empty() || record.signatures.is_empty())
114 {
115 return Err(Error::Invalid(
116 "original signatures and canonical records required",
117 ));
118 }
119 Ok(())
120 }
121}
122
123#[derive(Debug, thiserror::Error)]
124pub enum Error {
125 #[error(transparent)]
126 Client(#[from] ClientError<transport::Error>),
127 #[error(transparent)]
128 Transport(#[from] transport::Error),
129 #[error("publication source: {0}")]
130 Source(#[from] std::io::Error),
131 #[error("invalid publication: {0}")]
132 Invalid(&'static str),
133 #[error(
134 "publication {limit_name} limit is {limit}, required {actual} for {operations} operations"
135 )]
136 OriginalBudgetExceeded {
137 operations: usize,
138 limit_name: &'static str,
139 limit: usize,
140 actual: usize,
141 },
142 #[error(
143 "publication operation {operation} requires {bytes} encoded bytes; batch limit {batch_limit} bytes, frame limit {frame_limit} bytes"
144 )]
145 OriginalOperationTooLarge {
146 operation: usize,
147 bytes: usize,
148 batch_limit: usize,
149 frame_limit: usize,
150 },
151}
152
153impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
154 pub async fn publish_content<R: AsyncRead + Unpin + Send>(
159 &self,
160 opening: &PublishContentClientFrame,
161 originals: &PublicationOriginals,
162 mut artifacts: [R; 2],
163 ) -> Result<PublicationReceipt, Error> {
164 for batch in &originals.operations {
166 crate::hybrid::operations(batch).map_err(Error::Invalid)?;
167 }
168 let originals =
171 batching::bounded_originals(originals, &opening.client_operation_id, 512 * 1024)?;
172 originals.validate_bounds()?;
173 let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
174 return Err(Error::Invalid("Open required"));
175 };
176 crate::hybrid::publish_open(open).map_err(Error::Invalid)?;
177 if open.protocol.is_some() {
178 api::import_authority::require_hybrid_peer(self.description.protocol.as_ref())
179 .map_err(|_| Error::Invalid("peer does not support HYBRID publication"))?;
180 }
181 if originals.operations.iter().any(|batch| {
182 (batch.import_authority.is_some() && batch.import_authority != open.import_authority)
183 || (batch.native_authority.is_some()
184 && batch.native_authority != open.native_authority)
185 }) {
186 return Err(Error::Invalid(
187 "publication proof differs from negotiated opening",
188 ));
189 }
190 if open.destination != self.description.endpoint
191 || open.thread.is_none()
192 || open.revision.is_none()
193 || (!open.sharing_policy_version.is_empty() && open.sharing_policy_version.len() != 32)
194 || open.packs.len() != 2
195 || open.packs[0].kind != pack_extent::Kind::NativePack as i32
196 || open.packs[1].kind != pack_extent::Kind::NativeIndex as i32
197 {
198 return Err(Error::Invalid(
199 "endpoint, capture and ordered native artifacts required",
200 ));
201 }
202 let inventory = inventory_digest(&open.packs)?;
203 let mut logical = opening.clone();
204 if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
205 open.checkpoint = None;
206 }
207 let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
208 let expected = TransferCheckpoint {
209 transfer_id: digest[..16].to_vec(),
210 plan_digest: digest.to_vec(),
211 resume_token: Vec::new(),
212 committed_bytes: 0,
213 };
214 let (mut sender, mut responses) = self
215 .api
216 .exchange::<rpc::SyncServicePublishContent>(opening)
217 .await?;
218 let first = responses
219 .next()
220 .await?
221 .ok_or(Error::Invalid("publication ended before admission"))?;
222 let ready = match first.body {
223 Some(publish_content_server_frame::Body::Receipt(receipt)) => {
224 return validate_receipt(receipt, opening, &inventory);
225 }
226 Some(publish_content_server_frame::Body::Ready(ready)) => ready,
227 _ => return Err(Error::Invalid("Ready or replay receipt required")),
228 };
229 crate::hybrid::transfer_ready(&ready).map_err(Error::Invalid)?;
230 crate::hybrid::negotiated(open.protocol.as_ref(), ready.protocol.as_ref())
231 .map_err(Error::Invalid)?;
232 if ready.endpoint != open.destination
233 || ready.thread != open.thread
234 || ready.current != open.revision
235 || ready.checkpoint.as_ref() != Some(&expected)
236 {
237 return Err(Error::Invalid("admission differs from publication plan"));
238 }
239 let budget = ready
240 .budget
241 .ok_or(Error::Invalid("publication frame budget required"))?;
242 let frame_limit = budget.max_frame_bytes as usize;
243 if !(1024..=512 * 1024).contains(&frame_limit) {
244 return Err(Error::Invalid("unsupported publication frame budget"));
245 }
246 let originals =
247 batching::bounded_originals(&originals, &opening.client_operation_id, frame_limit)?;
248 originals.validate_bounds()?;
249 let upload = async {
250 for body in originals
251 .geneses
252 .iter()
253 .cloned()
254 .map(publish_content_client_frame::Body::ThreadGenesis)
255 .chain(
256 originals
257 .operations
258 .iter()
259 .cloned()
260 .map(publish_content_client_frame::Body::Operations),
261 )
262 {
263 let frame = PublishContentClientFrame {
264 client_operation_id: opening.client_operation_id.clone(),
265 body: Some(body),
266 };
267 if frame.encoded_len() > frame_limit {
268 return Err(Error::Invalid(
269 "original exceeds negotiated publication frame budget",
270 ));
271 }
272 sender.send(&frame).await?;
273 }
274 for (artifact, planned) in artifacts.iter_mut().zip(&open.packs) {
275 let mut offset = 0;
276 let mut digest = blake3::Hasher::new();
277 let mut buffer = vec![0; frame_limit / 2];
278 while offset < planned.length {
279 let length = (planned.length - offset).min(buffer.len() as u64) as usize;
280 let read = artifact.read(&mut buffer[..length]).await?;
281 if read == 0 {
282 return Err(Error::Invalid("artifact ended before declared length"));
283 }
284 let data = &buffer[..read];
285 digest.update(data);
286 let chunk = PublishContentClientFrame {
287 client_operation_id: opening.client_operation_id.clone(),
288 body: Some(publish_content_client_frame::Body::Pack(PackChunk {
289 extent: Some(PackExtent {
290 pack: planned.pack.clone(),
291 kind: planned.kind,
292 offset,
293 length: read as u64,
294 extent_digest: Some(ObjectAddress {
295 algorithm: "blake3".into(),
296 digest: blake3::hash(data).as_bytes().to_vec(),
297 }),
298 }),
299 data: data.to_vec(),
300 })),
301 };
302 if chunk.encoded_len() > frame_limit {
303 return Err(Error::Invalid("chunk exceeds negotiated frame budget"));
304 }
305 sender.send(&chunk).await?;
306 offset += read as u64;
307 }
308 if artifact.read(&mut buffer[..1]).await? != 0
309 || planned
310 .pack
311 .as_ref()
312 .is_none_or(|address| address.digest != digest.finalize().as_bytes())
313 {
314 return Err(Error::Invalid(
315 "artifact differs from declared length or digest",
316 ));
317 }
318 }
319 sender
320 .send(&PublishContentClientFrame {
321 client_operation_id: opening.client_operation_id.clone(),
322 body: Some(publish_content_client_frame::Body::Finish(
323 PublishContentFinish {
324 checkpoint: Some(expected.clone()),
325 },
326 )),
327 })
328 .await?;
329 sender.finish().await?;
330 Ok::<_, Error>(())
331 };
332 let receive = async {
333 while let Some(frame) = responses.next().await? {
334 match frame.body {
335 Some(publish_content_server_frame::Body::Checkpoint(checkpoint))
336 if checkpoint == expected => {}
337 Some(publish_content_server_frame::Body::Receipt(receipt)) => {
338 return validate_receipt(receipt, opening, &inventory);
339 }
340 _ => return Err(Error::Invalid("unexpected publication response")),
341 }
342 }
343 Err(Error::Invalid(
344 "publication ended without a durable receipt",
345 ))
346 };
347 let (_, receipt) = tokio::try_join!(upload, receive)?;
348 Ok(receipt)
349 }
350}
351
352fn validate_receipt(
353 receipt: PublicationReceipt,
354 opening: &PublishContentClientFrame,
355 inventory: &[u8; 32],
356) -> Result<PublicationReceipt, Error> {
357 let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
358 return Err(Error::Invalid("Open required"));
359 };
360 crate::hybrid::publication_receipt(&receipt).map_err(Error::Invalid)?;
361 if receipt.import_authority.is_some() || receipt.native_authority.is_some() {
362 api::import_authority::require_hybrid_peer(open.protocol.as_ref())
363 .map_err(|_| Error::Invalid("HYBRID receipt requires negotiated publication"))?;
364 }
365 if receipt.client_operation_id != opening.client_operation_id
366 || receipt.destination != open.destination
367 || receipt.thread != open.thread
368 || receipt.revision != open.revision
369 || receipt.sharing_policy_version.len() != 32
370 || (!open.sharing_policy_version.is_empty()
371 && receipt.sharing_policy_version != open.sharing_policy_version)
372 {
373 return Err(Error::Invalid("receipt differs from requested publication"));
374 }
375 match &receipt.outcome {
376 Some(publication_receipt::Outcome::Accepted(_))
377 if receipt.accepted_inventory.as_ref().is_some_and(|address| {
378 address.algorithm == "blake3" && address.digest == inventory
379 }) =>
380 {
381 Ok(receipt)
382 }
383 Some(publication_receipt::Outcome::Rejected(error)) => {
384 Err(transport::Error::Remote(error.clone().into()).into())
385 }
386 _ => Err(Error::Invalid(
387 "receipt did not accept the complete inventory",
388 )),
389 }
390}
391
392pub fn inventory_digest(packs: &[PackExtent]) -> Result<[u8; 32], Error> {
395 if packs.len() != 2
396 || packs[0].kind != pack_extent::Kind::NativePack as i32
397 || packs[1].kind != pack_extent::Kind::NativeIndex as i32
398 {
399 return Err(Error::Invalid("ordered native pack and index required"));
400 }
401 let mut inventory = Vec::new();
402 let mut total = 0_u64;
403 for extent in packs {
404 let address = extent
405 .pack
406 .as_ref()
407 .ok_or(Error::Invalid("artifact address required"))?;
408 total = total
409 .checked_add(extent.length)
410 .ok_or(Error::Invalid("artifact length overflow"))?;
411 if address.algorithm != "blake3"
412 || address.digest.len() != 32
413 || extent.length == 0
414 || extent.offset != 0
415 || extent.extent_digest.as_ref() != Some(address)
416 || total > 256 * 1024 * 1024
417 {
418 return Err(Error::Invalid("complete bounded BLAKE3 artifacts required"));
419 }
420 extent
421 .encode_length_delimited(&mut inventory)
422 .map_err(|_| Error::Invalid("inventory encoding failed"))?;
423 }
424 Ok(typed_digest("thread-source-inventory-v1", &inventory))
425}
426
427fn typed_digest(kind: &str, bytes: &[u8]) -> [u8; 32] {
428 let mut hasher = blake3::Hasher::new();
429 hasher.update(kind.as_bytes());
430 hasher.update(&(bytes.len() as u64).to_le_bytes());
431 hasher.update(&[0]);
432 hasher.update(bytes);
433 *hasher.finalize().as_bytes()
434}
435
436#[cfg(test)]
437mod tests {
438 use api::v2::{
439 MethodDescriptor,
440 client::{MessageReader, MessageWriter},
441 };
442 use tokio::sync::mpsc;
443
444 use super::*;
445
446 pub(super) struct Reader(mpsc::Receiver<Vec<u8>>);
447 pub(super) struct Writer(Option<mpsc::Sender<Vec<u8>>>);
448 impl MessageReader for Reader {
449 type Error = transport::Error;
450 async fn next(&mut self) -> Result<Option<Vec<u8>>, Self::Error> {
451 Ok(self.0.recv().await)
452 }
453 fn cancel(&mut self) {
454 self.0.close();
455 }
456 }
457 impl MessageWriter for Writer {
458 type Error = transport::Error;
459 async fn send(&mut self, bytes: Vec<u8>) -> Result<(), Self::Error> {
460 self.0
461 .as_ref()
462 .ok_or(transport::Error::Protocol("closed"))?
463 .send(bytes)
464 .await
465 .map_err(|_| transport::Error::Protocol("peer closed"))
466 }
467 async fn finish(&mut self) -> Result<(), Self::Error> {
468 self.0.take();
469 Ok(())
470 }
471 fn abort(&mut self) {
472 self.0.take();
473 }
474 }
475 pub(super) struct Peer {
476 wrong_receipt: bool,
477 original_count: usize,
478 }
479 impl RpcTransport for Peer {
480 type Error = transport::Error;
481 type Reader = Reader;
482 type Writer = Writer;
483 async fn unary(
484 &self,
485 _: &'static MethodDescriptor,
486 _: Vec<u8>,
487 ) -> Result<Vec<u8>, Self::Error> {
488 unreachable!("test exercises only publication");
489 }
490 async fn observe(
491 &self,
492 _: &'static MethodDescriptor,
493 _: Vec<u8>,
494 ) -> Result<Reader, Self::Error> {
495 unreachable!("test exercises only publication");
496 }
497 async fn exchange(
498 &self,
499 _: &'static MethodDescriptor,
500 bytes: Vec<u8>,
501 ) -> Result<(Writer, Reader), Self::Error> {
502 let opening = PublishContentClientFrame::decode(bytes.as_slice())?;
503 let Some(publish_content_client_frame::Body::Open(open)) = opening.body.clone() else {
504 panic!("Open");
505 };
506 let mut logical = opening.clone();
507 if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
508 open.checkpoint = None;
509 }
510 let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
511 let checkpoint = TransferCheckpoint {
512 transfer_id: digest[..16].to_vec(),
513 plan_digest: digest.to_vec(),
514 committed_bytes: 0,
515 resume_token: vec![],
516 };
517 let (tx, mut incoming) = mpsc::channel::<Vec<u8>>(1);
518 let (outgoing, rx) = mpsc::channel(1);
519 let wrong_receipt = self.wrong_receipt;
520 let original_count = self.original_count;
521 tokio::spawn(async move {
522 let send = |body| {
523 let outgoing = &outgoing;
524 async move {
525 outgoing
526 .send(PublishContentServerFrame { body: Some(body) }.encode_to_vec())
527 .await
528 }
529 };
530 send(publish_content_server_frame::Body::Ready(TransferReady {
531 endpoint: open.destination.clone(),
532 thread: open.thread.clone(),
533 current: open.revision.clone(),
534 checkpoint: Some(checkpoint.clone()),
535 protocol: open.protocol.clone(),
536 import_authority: open.import_authority.clone(),
537 budget: Some(ReadBudget {
538 max_frame_bytes: if open.import_authority.is_some() {
539 256 * 1024
540 } else {
541 2048
542 },
543 ..Default::default()
544 }),
545 ..Default::default()
546 }))
547 .await
548 .expect("Ready");
549 let mut lengths = [0_u64; 2];
550 let mut original_counts = [0usize; 2];
551 while let Some(bytes) = incoming.recv().await {
552 assert!(
553 bytes.len()
554 <= if open.import_authority.is_some() {
555 256 * 1024
556 } else {
557 2048
558 },
559 "negotiated frame size"
560 );
561 let frame = PublishContentClientFrame::decode(bytes.as_slice()).expect("frame");
562 assert_eq!(frame.client_operation_id, opening.client_operation_id);
563 match frame.body {
564 Some(publish_content_client_frame::Body::ThreadGenesis(genesis)) => {
565 assert_eq!(lengths, [0, 0]);
566 assert!(genesis.genesis.is_some());
567 original_counts[0] += 1;
568 send(publish_content_server_frame::Body::Checkpoint(
569 checkpoint.clone(),
570 ))
571 .await
572 .expect("original checkpoint");
573 }
574 Some(publish_content_client_frame::Body::Operations(batch)) => {
575 assert_eq!(lengths, [0, 0]);
576 assert!((1..=128).contains(&batch.operations.len()));
577 assert!(!batch.operations[0].canonical_record.is_empty());
578 original_counts[1] += batch.operations.len();
579 send(publish_content_server_frame::Body::Checkpoint(
580 checkpoint.clone(),
581 ))
582 .await
583 .expect("original checkpoint");
584 }
585 Some(publish_content_client_frame::Body::Pack(chunk)) => {
586 assert_eq!(
587 original_counts,
588 [1, original_count],
589 "original proofs precede source artifacts"
590 );
591 let extent = chunk.extent.expect("extent");
592 let index = if extent.kind == pack_extent::Kind::NativePack as i32 {
593 0
594 } else {
595 1
596 };
597 assert_eq!(extent.offset, lengths[index]);
598 assert_eq!(extent.length, chunk.data.len() as u64);
599 assert_eq!(
600 extent.extent_digest.expect("chunk digest").digest,
601 blake3::hash(&chunk.data).as_bytes()
602 );
603 lengths[index] += extent.length;
604 send(publish_content_server_frame::Body::Checkpoint(
608 checkpoint.clone(),
609 ))
610 .await
611 .expect("checkpoint");
612 }
613 Some(publish_content_client_frame::Body::Finish(finish)) => {
614 assert_eq!(finish.checkpoint, Some(checkpoint.clone()));
615 assert_eq!(lengths, [open.packs[0].length, open.packs[1].length]);
616 let mut inventory = Vec::new();
617 for extent in &open.packs {
618 extent
619 .encode_length_delimited(&mut inventory)
620 .expect("inventory");
621 }
622 let mut receipt = PublicationReceipt {
623 native_authority: None,
624 client_operation_id: opening.client_operation_id.clone(),
625 destination: open.destination.clone(),
626 thread: open.thread.clone(),
627 revision: open.revision.clone(),
628 sharing_policy_version: if open.sharing_policy_version.is_empty() {
629 vec![3; 32]
630 } else {
631 open.sharing_policy_version.clone()
632 },
633 accepted_inventory: Some(ObjectAddress {
634 algorithm: "blake3".into(),
635 digest: typed_digest("thread-source-inventory-v1", &inventory)
636 .to_vec(),
637 }),
638 outcome: Some(publication_receipt::Outcome::Accepted(
639 Applied::default(),
640 )),
641 import_authority: open.import_authority.clone(),
642 };
643 if wrong_receipt {
644 receipt.thread = None;
645 }
646 let _ =
647 send(publish_content_server_frame::Body::Receipt(receipt)).await;
648 break;
649 }
650 _ => panic!("unexpected frame"),
651 }
652 }
653 });
654 Ok((Writer(Some(tx)), Reader(rx)))
655 }
656 }
657
658 pub(super) fn fixture(
659 wrong_receipt: bool,
660 ) -> (
661 Remote<Peer>,
662 PublishContentClientFrame,
663 [std::io::Cursor<Vec<u8>>; 2],
664 ) {
665 let endpoint = EndpointRef {
666 public_key: vec![8; 32],
667 kind: EndpointKind::Weft as i32,
668 };
669 let remote = Remote {
670 api: api::v2::client::Client::new(
671 Peer {
672 wrong_receipt,
673 original_count: 1,
674 },
675 ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
676 ),
677 description: DescribeEndpointResponse {
678 endpoint: Some(endpoint.clone()),
679 ..Default::default()
680 },
681 };
682 let artifacts = [vec![1; 128 * 1024], vec![2; 16 * 1024]];
683 let packs = artifacts
684 .iter()
685 .zip([
686 pack_extent::Kind::NativePack,
687 pack_extent::Kind::NativeIndex,
688 ])
689 .map(|(data, kind)| {
690 let address = ObjectAddress {
691 algorithm: "blake3".into(),
692 digest: blake3::hash(data).as_bytes().to_vec(),
693 };
694 PackExtent {
695 pack: Some(address.clone()),
696 kind: kind as i32,
697 offset: 0,
698 length: data.len() as u64,
699 extent_digest: Some(address),
700 }
701 })
702 .collect();
703 let open = PublishContentClientFrame {
704 client_operation_id: "op-test".into(),
705 body: Some(publish_content_client_frame::Body::Open(
706 PublishContentOpen {
707 destination: Some(endpoint),
708 thread: Some(ThreadRef::default()),
709 revision: Some(RevisionRef::default()),
710 packs,
711 ..Default::default()
712 },
713 )),
714 };
715 (remote, open, artifacts.map(std::io::Cursor::new))
716 }
717 fn originals() -> PublicationOriginals {
720 let record = SignedRecord {
721 format: "transport-fixture".into(),
722 canonical_record: vec![1],
723 signatures: vec![RecordSignature {
724 public_key: vec![2; 32],
725 signature: vec![3; 64],
726 }],
727 };
728 PublicationOriginals {
729 geneses: vec![ThreadGenesisRecord {
730 boundary_acceptances: Vec::new(),
731 genesis: Some(record.clone()),
732 ..Default::default()
733 }],
734 operations: vec![ReplicationOperations {
735 native_authority: None,
736 boundary_acceptances: Vec::new(),
737 operations: vec![record],
738 authority_admissions: vec![],
739 import_authority: None,
740 }],
741 }
742 }
743
744 #[test]
745 fn publication_originals_require_bounded_complete_metadata() {
746 let valid = originals();
747 assert!(valid.validate_bounds().is_ok());
748 let mut missing = valid.clone();
749 missing.operations.clear();
750 assert!(matches!(
751 missing.validate_bounds(),
752 Err(Error::Invalid(
753 "bounded original genesis and source operations required"
754 ))
755 ));
756 let mut unsigned = valid.clone();
757 unsigned.operations[0].operations[0].signatures.clear();
758 assert!(matches!(
759 unsigned.validate_bounds(),
760 Err(Error::Invalid(
761 "original signatures and canonical records required"
762 ))
763 ));
764 let mut oversized = valid.clone();
765 oversized.operations[0].operations[0].canonical_record = vec![1; 256 * 1024];
766 assert!(matches!(
767 oversized.validate_bounds(),
768 Err(Error::Invalid(
769 "publication original metadata budget exceeded"
770 ))
771 ));
772 let mut too_many = valid;
773 too_many.geneses = vec![too_many.geneses[0].clone(); 129];
774 assert!(matches!(
775 too_many.validate_bounds(),
776 Err(Error::Invalid(
777 "bounded original genesis and source operations required"
778 ))
779 ));
780 }
781
782 #[tokio::test]
783 async fn publication_drains_checkpoints_while_uploading_under_backpressure() {
784 let (remote, open, artifacts) = fixture(false);
785 let receipt = tokio::time::timeout(
786 std::time::Duration::from_secs(2),
787 remote.publish_content(&open, &originals(), artifacts),
788 )
789 .await
790 .expect("must not deadlock")
791 .expect("publication");
792 assert_eq!(receipt.client_operation_id, "op-test");
793 }
794 #[tokio::test]
795 async fn publication_rebatches_after_ready_before_upload() {
796 let (mut remote, open, artifacts) = fixture(false);
797 remote.api = api::v2::client::Client::new(
798 Peer {
799 wrong_receipt: false,
800 original_count: 140,
801 },
802 ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
803 );
804 let mut originals = originals();
805 let mut record = originals.operations[0].operations[0].clone();
806 record.canonical_record = vec![1; 300];
807 originals.operations[0].operations = vec![record; 140];
808 let receipt = tokio::time::timeout(
809 std::time::Duration::from_secs(5),
810 remote.publish_content(&open, &originals, artifacts),
811 )
812 .await
813 .expect("must not deadlock")
814 .expect("negotiated bounded publication");
815 assert_eq!(receipt.client_operation_id, "op-test");
816 }
817
818 #[cfg(feature = "native")]
819 #[tokio::test]
820 async fn hosted_source_publication_retains_complete_history_and_original_signatures() {
821 use objects::store::{FsStore, ObjectStore};
822 let scratch = tempfile::tempdir().expect("source");
823 let (staged, _, pinned) = crate::fetch::hosted::tests::source(scratch.path(), false);
824 let bundle = staged.import_authority().expect("complete history").clone();
825 let store = FsStore::new(scratch.path().join("publication-store"));
826 store.init().expect("store");
827 let paths = staged.artifact_paths();
828 store
829 .install_pack_streaming(&paths[0], &paths[1])
830 .expect("isolated original pack");
831 let pack = SourcePack::prepare(
832 &store,
833 staged.state(),
834 scratch.path(),
835 SourceBudget {
836 max_objects: 16,
837 max_decoded_bytes: 1024 * 1024,
838 },
839 )
840 .expect("selected closure");
841 let operations = crate::authority_admission::batches(
842 staged.operations().iter().cloned().map(|original| {
843 crate::replication::store::ReceivedOperation {
844 native_authority: None,
845 original,
846 authority_admission: None,
847 import_authority: Some(std::sync::Arc::new(bundle.clone())),
848 }
849 }),
850 128 * 1024,
851 128,
852 )
853 .expect("bounded carrier")
854 .collect::<Result<Vec<_>, _>>()
855 .expect("originals");
856 let originals = PublicationOriginals {
857 geneses: vec![
858 staged
859 .ready()
860 .thread_genesis
861 .clone()
862 .expect("selected genesis"),
863 ],
864 operations,
865 };
866 let selected = staged.ready().thread.clone().expect("Thread");
867 let (mut remote, _, _) = fixture(false);
868 let options = || PublicationOptions {
869 client_operation_id: "2a2a2a2a-2a2a-2a2a-2a2a-2a2a2a2a2a2a".into(),
870 source: EndpointRef {
871 public_key: vec![2; 32],
872 kind: EndpointKind::Device as i32,
873 },
874 sharing_policy_version: vec![],
875 checkpoint: None,
876 };
877 assert!(
878 remote
879 .thread(selected.clone())
880 .publish_source(&pack, &originals, options())
881 .await
882 .is_err(),
883 "old peer cannot strip history"
884 );
885 remote.description.protocol = Some(crate::hybrid::protocol());
886 let receipt = remote
887 .thread(selected)
888 .publish_source(&pack, &originals, options())
889 .await
890 .expect("capable publication");
891 assert_eq!(receipt.import_authority.as_ref(), Some(&bundle));
892 let history = crate::hybrid::authority::AcceptedHistory::from_selected_spool(
893 &bundle,
894 &pinned,
895 1350,
896 heddleco_capability_verifier::VerificationLimits::new(30 * 24 * 60 * 60)
897 .expect("limits"),
898 )
899 .expect("independent owner");
900 let prepared = remote
901 .thread(staged.ready().thread.clone().expect("Thread"))
902 .prepare_publication(
903 &pack,
904 originals.clone(),
905 PublicationOptions {
906 client_operation_id: "2b2b2b2b-2b2b-2b2b-2b2b-2b2b2b2b2b2b".into(),
907 source: EndpointRef {
908 public_key: vec![2; 32],
909 kind: EndpointKind::Device as i32,
910 },
911 sharing_policy_version: vec![],
912 checkpoint: None,
913 },
914 heddle_object_model::object::ContentHash::from_bytes(*history.genesis()),
915 )
916 .expect("prepared public originals");
917 let Some(publish_content_client_frame::Body::Open(open)) = &prepared.opening().body else {
918 panic!("Open")
919 };
920 assert_eq!(open.import_authority.as_ref(), Some(&bundle));
921 let received = tempfile::tempdir().expect("received artifacts");
922 for (mut file, name) in pack
923 .open_artifacts()
924 .await
925 .expect("original artifacts")
926 .into_iter()
927 .zip(["source.pack", "source.idx"])
928 {
929 let mut output = tokio::fs::File::create(received.path().join(name))
930 .await
931 .expect("received file");
932 tokio::io::copy(&mut file, &mut output)
933 .await
934 .expect("exact uploaded bytes");
935 }
936 let authority = crate::hybrid::authority::SelectedAuthority::new(
937 history,
938 bundle.clone(),
939 |_: &ImportPublicProofBundleV1,
940 _: i64,
941 _: &repo::thread_replication::hosted_trust::TrustTransaction<'_>| Ok(()),
942 );
943 let fixture: serde_json::Value =
944 serde_json::from_str(include_str!("../tests/fixtures/hybrid-alpha33.json"))
945 .expect("tagged vectors");
946 let carriers = repo::thread_replication::delegated_import::authenticate_import_carriers(
947 &bundle,
948 &authority,
949 &api::import_authority::ImportWitnessRootPin {
950 authority: "https://weft.example.test".into(),
951 root_id: "descriptor-root-1".into(),
952 public_key: hex::decode(
953 fixture["keys"]["root"]["public_key_hex"]
954 .as_str()
955 .expect("root"),
956 )
957 .expect("root key"),
958 epoch: 1,
959 },
960 1_350_000,
961 &[],
962 &[],
963 |_| Ok(()),
964 )
965 .expect("independently authenticated import carriers");
966 let validated = validate_source_artifacts_with_import_carriers(
967 received,
968 open,
969 prepared.originals().clone(),
970 carriers,
971 )
972 .expect("signed originals and actual closure");
973 assert_eq!(validated.import_authority(), Some(&bundle));
974 let received = validated
975 .into_hosted_source(staged.ready().clone())
976 .expect("retain public history through staging");
977 assert_eq!(received.import_authority(), Some(&bundle));
978 assert_eq!(received.operations(), staged.operations());
979
980 assert_eq!(
981 originals.operations[0].operations[0].canonical_record,
982 staged.operations()[0].canonical
983 );
984 assert_eq!(
985 originals.operations[0].operations[0].signatures[0].signature,
986 staged.operations()[0].signature
987 );
988 assert_eq!(
989 bundle.original_geneses.len(),
990 2,
991 "complete sibling public history survives selected source publication"
992 );
993 }
994
995 #[tokio::test]
996 async fn publication_refuses_originals_carrying_hybrid_import_authority() {
997 let (remote, open, artifacts) = fixture(false);
998 let mut originals = originals();
999 originals.operations[0].import_authority = Some(ImportPublicProofBundleV1::default());
1000 let error = remote
1001 .publish_content(&open, &originals, artifacts)
1002 .await
1003 .expect_err("HYBRID originals must not be relayed with the bundle dropped");
1004 assert!(
1005 matches!(error, Error::Invalid(message) if message.contains("api#307")),
1006 "{error}"
1007 );
1008 }
1009
1010 #[tokio::test]
1011 async fn publication_refuses_a_receipt_for_another_scope() {
1012 let (remote, open, artifacts) = fixture(true);
1013 let error = remote
1014 .publish_content(&open, &originals(), artifacts)
1015 .await
1016 .expect_err("mismatched scope");
1017 assert!(matches!(
1018 error,
1019 Error::Invalid("receipt differs from requested publication")
1020 ));
1021 }
1022 #[test]
1023 fn publication_policy_is_optional_cas_and_receipt_reports_actual_frontier() {
1024 let (_, mut opening, _) = fixture(false);
1025 let inventory = [7; 32];
1026 let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
1027 panic!("opening")
1028 };
1029 let receipt = PublicationReceipt {
1030 native_authority: None,
1031 client_operation_id: opening.client_operation_id.clone(),
1032 destination: open.destination.clone(),
1033 thread: open.thread.clone(),
1034 revision: open.revision.clone(),
1035 sharing_policy_version: vec![3; 32],
1036 accepted_inventory: Some(ObjectAddress {
1037 algorithm: "blake3".into(),
1038 digest: inventory.to_vec(),
1039 }),
1040 outcome: Some(publication_receipt::Outcome::Accepted(Applied::default())),
1041 import_authority: None,
1042 };
1043 validate_receipt(receipt.clone(), &opening, &inventory)
1044 .expect("one-shot upload returns actual policy without CAS");
1045 let mut hybrid = receipt.clone();
1046 hybrid.import_authority = Some(ImportPublicProofBundleV1::default());
1047 assert!(
1048 matches!(
1049 validate_receipt(hybrid, &opening, &inventory),
1050 Err(Error::Invalid(message)) if message.contains("api#307")
1051 ),
1052 "a receipt carrying HYBRID import authority must be refused, never ignored"
1053 );
1054 let mut absent = receipt.clone();
1055 absent.sharing_policy_version.clear();
1056 assert!(
1057 validate_receipt(absent, &opening, &inventory).is_err(),
1058 "receipt must report actual policy frontier"
1059 );
1060 let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
1061 panic!("opening")
1062 };
1063 open.sharing_policy_version = vec![4; 32];
1064 assert!(
1065 validate_receipt(receipt.clone(), &opening, &inventory).is_err(),
1066 "explicit policy CAS cannot silently accept another version"
1067 );
1068 let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
1069 panic!("opening")
1070 };
1071 open.sharing_policy_version = vec![3; 32];
1072 validate_receipt(receipt, &opening, &inventory).expect("matching explicit policy version");
1073 }
1074
1075 #[cfg(feature = "replication")]
1076 #[test]
1077 fn publication_digests_match_the_native_typed_hash_format() {
1078 assert_eq!(
1079 typed_digest("thread-source-inventory-v1", b"bytes"),
1080 *heddle_object_model::object::ContentHash::compute_typed(
1081 "thread-source-inventory-v1",
1082 b"bytes"
1083 )
1084 .as_bytes()
1085 );
1086 }
1087
1088 #[cfg(feature = "source-transfer")]
1089 #[tokio::test]
1090 async fn thread_publication_prepares_only_selected_source_and_binds_its_revision() {
1091 use objects::{
1092 object::{Attribution, Blob, Principal, State, Tree, TreeEntry},
1093 store::{FsStore, ObjectStore},
1094 };
1095 let root = tempfile::tempdir().expect("local source scratch");
1096 let store = FsStore::new(root.path().join("objects"));
1097 store.init().expect("store");
1098 let blob = Blob::new(vec![1; 128 * 1024]);
1099 store.put_blob(&blob).expect("source");
1100 let tree = Tree::from_entries(vec![
1101 TreeEntry::file("source.rs", blob.hash(), false).expect("entry"),
1102 ]);
1103 store.put_tree(&tree).expect("tree");
1104 let state = State::new_snapshot(
1105 tree.hash(),
1106 vec![],
1107 Attribution::human(Principal::new("user", "user@example.test")),
1108 );
1109 let prepared = SourcePack::prepare(
1110 &store,
1111 &state,
1112 root.path(),
1113 SourceBudget {
1114 max_objects: 16,
1115 max_decoded_bytes: 256 * 1024,
1116 },
1117 )
1118 .expect("prepare exact source");
1119 let (remote, _, _) = fixture(false);
1120 let thread = ThreadRef {
1121 spool: Some(SpoolRef {
1122 id: "spool-test".into(),
1123 }),
1124 id: Some(ThreadId { value: vec![1; 32] }),
1125 };
1126 let receipt = remote
1127 .thread(thread.clone())
1128 .publish_source(
1129 &prepared,
1130 &originals(),
1131 PublicationOptions {
1132 client_operation_id: "source-upload".into(),
1133 source: EndpointRef {
1134 public_key: vec![2; 32],
1135 kind: EndpointKind::Device as i32,
1136 },
1137 sharing_policy_version: vec![],
1138 checkpoint: None,
1139 },
1140 )
1141 .await
1142 .expect("one Thread-bound publication");
1143 assert_eq!(receipt.thread, Some(thread.clone()));
1144 assert_eq!(
1145 receipt.revision,
1146 Some(RevisionRef {
1147 spool: thread.spool,
1148 revision: Some(revision_ref::Revision::State(
1149 api::heddle::api::common::StateId {
1150 value: state.id().as_bytes().to_vec()
1151 }
1152 ))
1153 })
1154 );
1155 }
1156}