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};
29
30#[derive(Clone)]
34pub struct PublicationOriginals {
35 pub geneses: Vec<ThreadGenesisRecord>,
36 pub operations: Vec<ReplicationOperations>,
37}
38impl PublicationOriginals {
39 fn validate_bounds(&self) -> Result<(), Error> {
40 if self.geneses.is_empty()
41 || self.geneses.len() > 128
42 || self.operations.is_empty()
43 || self.operations.len() > 10_000
44 || self
45 .operations
46 .iter()
47 .map(|batch| batch.operations.len())
48 .sum::<usize>()
49 > 10_000
50 || self.operations.iter().any(|batch| {
51 batch.operations.is_empty()
52 || batch.authority_admissions.len() > batch.operations.len()
53 })
54 {
55 return Err(Error::Invalid(
56 "bounded original genesis and source operations required",
57 ));
58 }
59 let mut acceptances = std::collections::BTreeSet::new();
60 for record in self
61 .geneses
62 .iter()
63 .flat_map(|value| &value.boundary_acceptances)
64 .chain(
65 self.operations
66 .iter()
67 .flat_map(|value| &value.boundary_acceptances),
68 )
69 {
70 if record.canonical_record.len() > 96 * 1024
71 || record.signatures.len() != 1
72 || record.signatures[0].signature.len() != 64
73 || record.signatures[0].public_key.len() != 32
74 {
75 return Err(Error::Invalid("boundary evidence shape exceeds bounds"));
76 }
77 acceptances.insert(record.canonical_record.as_slice());
78 if acceptances.len() > 128 {
79 return Err(Error::Invalid("boundary acceptance count exceeded"));
80 }
81 }
82 let mut bytes = 0usize;
83 for length in self
84 .geneses
85 .iter()
86 .map(Message::encoded_len)
87 .chain(self.operations.iter().map(Message::encoded_len))
88 {
89 bytes = bytes
90 .checked_add(length)
91 .ok_or(Error::Invalid("publication original size overflow"))?;
92 if length > 256 * 1024 || bytes > 16 * 1024 * 1024 {
93 return Err(Error::Invalid(
94 "publication original metadata budget exceeded",
95 ));
96 }
97 }
98 if self.geneses.iter().any(|genesis| {
99 genesis.genesis.as_ref().is_none_or(|record| {
100 record.canonical_record.is_empty() || record.signatures.is_empty()
101 })
102 }) || self
103 .operations
104 .iter()
105 .flat_map(|batch| {
106 batch
107 .operations
108 .iter()
109 .chain(&batch.authority_admissions)
110 .chain(&batch.boundary_acceptances)
111 })
112 .any(|record| record.canonical_record.is_empty() || record.signatures.is_empty())
113 {
114 return Err(Error::Invalid(
115 "original signatures and canonical records required",
116 ));
117 }
118 Ok(())
119 }
120}
121
122#[derive(Debug, thiserror::Error)]
123pub enum Error {
124 #[error(transparent)]
125 Client(#[from] ClientError<transport::Error>),
126 #[error(transparent)]
127 Transport(#[from] transport::Error),
128 #[error("publication source: {0}")]
129 Source(#[from] std::io::Error),
130 #[error("invalid publication: {0}")]
131 Invalid(&'static str),
132 #[error(
133 "publication {limit_name} limit is {limit}, required {actual} for {operations} operations"
134 )]
135 OriginalBudgetExceeded {
136 operations: usize,
137 limit_name: &'static str,
138 limit: usize,
139 actual: usize,
140 },
141 #[error(
142 "publication operation {operation} requires {bytes} encoded bytes; batch limit {batch_limit} bytes, frame limit {frame_limit} bytes"
143 )]
144 OriginalOperationTooLarge {
145 operation: usize,
146 bytes: usize,
147 batch_limit: usize,
148 frame_limit: usize,
149 },
150}
151
152impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
153 pub async fn publish_content<R: AsyncRead + Unpin + Send>(
158 &self,
159 opening: &PublishContentClientFrame,
160 originals: &PublicationOriginals,
161 mut artifacts: [R; 2],
162 ) -> Result<PublicationReceipt, Error> {
163 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 if open.destination != self.description.endpoint
177 || open.thread.is_none()
178 || open.revision.is_none()
179 || (!open.sharing_policy_version.is_empty() && open.sharing_policy_version.len() != 32)
180 || open.packs.len() != 2
181 || open.packs[0].kind != pack_extent::Kind::NativePack as i32
182 || open.packs[1].kind != pack_extent::Kind::NativeIndex as i32
183 {
184 return Err(Error::Invalid(
185 "endpoint, capture and ordered native artifacts required",
186 ));
187 }
188 let inventory = inventory_digest(&open.packs)?;
189 let mut logical = opening.clone();
190 if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
191 open.checkpoint = None;
192 }
193 let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
194 let expected = TransferCheckpoint {
195 transfer_id: digest[..16].to_vec(),
196 plan_digest: digest.to_vec(),
197 resume_token: Vec::new(),
198 committed_bytes: 0,
199 };
200 let (mut sender, mut responses) = self
201 .api
202 .exchange::<rpc::SyncServicePublishContent>(opening)
203 .await?;
204 let first = responses
205 .next()
206 .await?
207 .ok_or(Error::Invalid("publication ended before admission"))?;
208 let ready = match first.body {
209 Some(publish_content_server_frame::Body::Receipt(receipt)) => {
210 return validate_receipt(receipt, opening, &inventory);
211 }
212 Some(publish_content_server_frame::Body::Ready(ready)) => ready,
213 _ => return Err(Error::Invalid("Ready or replay receipt required")),
214 };
215 crate::hybrid::transfer_ready(&ready).map_err(Error::Invalid)?;
216 if ready.endpoint != open.destination
217 || ready.thread != open.thread
218 || ready.current != open.revision
219 || ready.checkpoint.as_ref() != Some(&expected)
220 {
221 return Err(Error::Invalid("admission differs from publication plan"));
222 }
223 let budget = ready
224 .budget
225 .ok_or(Error::Invalid("publication frame budget required"))?;
226 let frame_limit = budget.max_frame_bytes as usize;
227 if !(1024..=512 * 1024).contains(&frame_limit) {
228 return Err(Error::Invalid("unsupported publication frame budget"));
229 }
230 let originals =
231 batching::bounded_originals(&originals, &opening.client_operation_id, frame_limit)?;
232 originals.validate_bounds()?;
233 let upload = async {
234 for body in originals
235 .geneses
236 .iter()
237 .cloned()
238 .map(publish_content_client_frame::Body::ThreadGenesis)
239 .chain(
240 originals
241 .operations
242 .iter()
243 .cloned()
244 .map(publish_content_client_frame::Body::Operations),
245 )
246 {
247 let frame = PublishContentClientFrame {
248 client_operation_id: opening.client_operation_id.clone(),
249 body: Some(body),
250 };
251 if frame.encoded_len() > frame_limit {
252 return Err(Error::Invalid(
253 "original exceeds negotiated publication frame budget",
254 ));
255 }
256 sender.send(&frame).await?;
257 }
258 for (artifact, planned) in artifacts.iter_mut().zip(&open.packs) {
259 let mut offset = 0;
260 let mut digest = blake3::Hasher::new();
261 let mut buffer = vec![0; frame_limit / 2];
262 while offset < planned.length {
263 let length = (planned.length - offset).min(buffer.len() as u64) as usize;
264 let read = artifact.read(&mut buffer[..length]).await?;
265 if read == 0 {
266 return Err(Error::Invalid("artifact ended before declared length"));
267 }
268 let data = &buffer[..read];
269 digest.update(data);
270 let chunk = PublishContentClientFrame {
271 client_operation_id: opening.client_operation_id.clone(),
272 body: Some(publish_content_client_frame::Body::Pack(PackChunk {
273 extent: Some(PackExtent {
274 pack: planned.pack.clone(),
275 kind: planned.kind,
276 offset,
277 length: read as u64,
278 extent_digest: Some(ObjectAddress {
279 algorithm: "blake3".into(),
280 digest: blake3::hash(data).as_bytes().to_vec(),
281 }),
282 }),
283 data: data.to_vec(),
284 })),
285 };
286 if chunk.encoded_len() > frame_limit {
287 return Err(Error::Invalid("chunk exceeds negotiated frame budget"));
288 }
289 sender.send(&chunk).await?;
290 offset += read as u64;
291 }
292 if artifact.read(&mut buffer[..1]).await? != 0
293 || planned
294 .pack
295 .as_ref()
296 .is_none_or(|address| address.digest != digest.finalize().as_bytes())
297 {
298 return Err(Error::Invalid(
299 "artifact differs from declared length or digest",
300 ));
301 }
302 }
303 sender
304 .send(&PublishContentClientFrame {
305 client_operation_id: opening.client_operation_id.clone(),
306 body: Some(publish_content_client_frame::Body::Finish(
307 PublishContentFinish {
308 checkpoint: Some(expected.clone()),
309 },
310 )),
311 })
312 .await?;
313 sender.finish().await?;
314 Ok::<_, Error>(())
315 };
316 let receive = async {
317 while let Some(frame) = responses.next().await? {
318 match frame.body {
319 Some(publish_content_server_frame::Body::Checkpoint(checkpoint))
320 if checkpoint == expected => {}
321 Some(publish_content_server_frame::Body::Receipt(receipt)) => {
322 return validate_receipt(receipt, opening, &inventory);
323 }
324 _ => return Err(Error::Invalid("unexpected publication response")),
325 }
326 }
327 Err(Error::Invalid(
328 "publication ended without a durable receipt",
329 ))
330 };
331 let (_, receipt) = tokio::try_join!(upload, receive)?;
332 Ok(receipt)
333 }
334}
335
336fn validate_receipt(
337 receipt: PublicationReceipt,
338 opening: &PublishContentClientFrame,
339 inventory: &[u8; 32],
340) -> Result<PublicationReceipt, Error> {
341 let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
342 return Err(Error::Invalid("Open required"));
343 };
344 crate::hybrid::publication_receipt(&receipt).map_err(Error::Invalid)?;
345 if receipt.client_operation_id != opening.client_operation_id
346 || receipt.destination != open.destination
347 || receipt.thread != open.thread
348 || receipt.revision != open.revision
349 || receipt.sharing_policy_version.len() != 32
350 || (!open.sharing_policy_version.is_empty()
351 && receipt.sharing_policy_version != open.sharing_policy_version)
352 {
353 return Err(Error::Invalid("receipt differs from requested publication"));
354 }
355 match &receipt.outcome {
356 Some(publication_receipt::Outcome::Accepted(_))
357 if receipt.accepted_inventory.as_ref().is_some_and(|address| {
358 address.algorithm == "blake3" && address.digest == inventory
359 }) =>
360 {
361 Ok(receipt)
362 }
363 Some(publication_receipt::Outcome::Rejected(error)) => {
364 Err(transport::Error::Remote(error.clone().into()).into())
365 }
366 _ => Err(Error::Invalid(
367 "receipt did not accept the complete inventory",
368 )),
369 }
370}
371
372pub fn inventory_digest(packs: &[PackExtent]) -> Result<[u8; 32], Error> {
375 if packs.len() != 2
376 || packs[0].kind != pack_extent::Kind::NativePack as i32
377 || packs[1].kind != pack_extent::Kind::NativeIndex as i32
378 {
379 return Err(Error::Invalid("ordered native pack and index required"));
380 }
381 let mut inventory = Vec::new();
382 let mut total = 0_u64;
383 for extent in packs {
384 let address = extent
385 .pack
386 .as_ref()
387 .ok_or(Error::Invalid("artifact address required"))?;
388 total = total
389 .checked_add(extent.length)
390 .ok_or(Error::Invalid("artifact length overflow"))?;
391 if address.algorithm != "blake3"
392 || address.digest.len() != 32
393 || extent.length == 0
394 || extent.offset != 0
395 || extent.extent_digest.as_ref() != Some(address)
396 || total > 256 * 1024 * 1024
397 {
398 return Err(Error::Invalid("complete bounded BLAKE3 artifacts required"));
399 }
400 extent
401 .encode_length_delimited(&mut inventory)
402 .map_err(|_| Error::Invalid("inventory encoding failed"))?;
403 }
404 Ok(typed_digest("thread-source-inventory-v1", &inventory))
405}
406
407fn typed_digest(kind: &str, bytes: &[u8]) -> [u8; 32] {
408 let mut hasher = blake3::Hasher::new();
409 hasher.update(kind.as_bytes());
410 hasher.update(&(bytes.len() as u64).to_le_bytes());
411 hasher.update(&[0]);
412 hasher.update(bytes);
413 *hasher.finalize().as_bytes()
414}
415
416#[cfg(test)]
417mod tests {
418 use api::v2::{
419 MethodDescriptor,
420 client::{MessageReader, MessageWriter},
421 };
422 use tokio::sync::mpsc;
423
424 use super::*;
425
426 pub(super) struct Reader(mpsc::Receiver<Vec<u8>>);
427 pub(super) struct Writer(Option<mpsc::Sender<Vec<u8>>>);
428 impl MessageReader for Reader {
429 type Error = transport::Error;
430 async fn next(&mut self) -> Result<Option<Vec<u8>>, Self::Error> {
431 Ok(self.0.recv().await)
432 }
433 fn cancel(&mut self) {
434 self.0.close();
435 }
436 }
437 impl MessageWriter for Writer {
438 type Error = transport::Error;
439 async fn send(&mut self, bytes: Vec<u8>) -> Result<(), Self::Error> {
440 self.0
441 .as_ref()
442 .ok_or(transport::Error::Protocol("closed"))?
443 .send(bytes)
444 .await
445 .map_err(|_| transport::Error::Protocol("peer closed"))
446 }
447 async fn finish(&mut self) -> Result<(), Self::Error> {
448 self.0.take();
449 Ok(())
450 }
451 fn abort(&mut self) {
452 self.0.take();
453 }
454 }
455 pub(super) struct Peer {
456 wrong_receipt: bool,
457 original_count: usize,
458 }
459 impl RpcTransport for Peer {
460 type Error = transport::Error;
461 type Reader = Reader;
462 type Writer = Writer;
463 async fn unary(
464 &self,
465 _: &'static MethodDescriptor,
466 _: Vec<u8>,
467 ) -> Result<Vec<u8>, Self::Error> {
468 unreachable!("test exercises only publication");
469 }
470 async fn observe(
471 &self,
472 _: &'static MethodDescriptor,
473 _: Vec<u8>,
474 ) -> Result<Reader, Self::Error> {
475 unreachable!("test exercises only publication");
476 }
477 async fn exchange(
478 &self,
479 _: &'static MethodDescriptor,
480 bytes: Vec<u8>,
481 ) -> Result<(Writer, Reader), Self::Error> {
482 let opening = PublishContentClientFrame::decode(bytes.as_slice())?;
483 let Some(publish_content_client_frame::Body::Open(open)) = opening.body.clone() else {
484 panic!("Open");
485 };
486 let mut logical = opening.clone();
487 if let Some(publish_content_client_frame::Body::Open(open)) = logical.body.as_mut() {
488 open.checkpoint = None;
489 }
490 let digest = typed_digest("thread-source-transfer-v1", &logical.encode_to_vec());
491 let checkpoint = TransferCheckpoint {
492 transfer_id: digest[..16].to_vec(),
493 plan_digest: digest.to_vec(),
494 committed_bytes: 0,
495 resume_token: vec![],
496 };
497 let (tx, mut incoming) = mpsc::channel::<Vec<u8>>(1);
498 let (outgoing, rx) = mpsc::channel(1);
499 let wrong_receipt = self.wrong_receipt;
500 let original_count = self.original_count;
501 tokio::spawn(async move {
502 let send = |body| {
503 let outgoing = &outgoing;
504 async move {
505 outgoing
506 .send(PublishContentServerFrame { body: Some(body) }.encode_to_vec())
507 .await
508 }
509 };
510 send(publish_content_server_frame::Body::Ready(TransferReady {
511 endpoint: open.destination.clone(),
512 thread: open.thread.clone(),
513 current: open.revision.clone(),
514 checkpoint: Some(checkpoint.clone()),
515 budget: Some(ReadBudget {
516 max_frame_bytes: 2048,
517 ..Default::default()
518 }),
519 ..Default::default()
520 }))
521 .await
522 .expect("Ready");
523 let mut lengths = [0_u64; 2];
524 let mut original_counts = [0usize; 2];
525 while let Some(bytes) = incoming.recv().await {
526 assert!(bytes.len() <= 2048, "negotiated frame size");
527 let frame = PublishContentClientFrame::decode(bytes.as_slice()).expect("frame");
528 assert_eq!(frame.client_operation_id, opening.client_operation_id);
529 match frame.body {
530 Some(publish_content_client_frame::Body::ThreadGenesis(genesis)) => {
531 assert_eq!(lengths, [0, 0]);
532 assert!(genesis.genesis.is_some());
533 original_counts[0] += 1;
534 send(publish_content_server_frame::Body::Checkpoint(
535 checkpoint.clone(),
536 ))
537 .await
538 .expect("original checkpoint");
539 }
540 Some(publish_content_client_frame::Body::Operations(batch)) => {
541 assert_eq!(lengths, [0, 0]);
542 assert!((1..=128).contains(&batch.operations.len()));
543 assert!(!batch.operations[0].canonical_record.is_empty());
544 original_counts[1] += batch.operations.len();
545 send(publish_content_server_frame::Body::Checkpoint(
546 checkpoint.clone(),
547 ))
548 .await
549 .expect("original checkpoint");
550 }
551 Some(publish_content_client_frame::Body::Pack(chunk)) => {
552 assert_eq!(
553 original_counts,
554 [1, original_count],
555 "original proofs precede source artifacts"
556 );
557 let extent = chunk.extent.expect("extent");
558 let index = if extent.kind == pack_extent::Kind::NativePack as i32 {
559 0
560 } else {
561 1
562 };
563 assert_eq!(extent.offset, lengths[index]);
564 assert_eq!(extent.length, chunk.data.len() as u64);
565 assert_eq!(
566 extent.extent_digest.expect("chunk digest").digest,
567 blake3::hash(&chunk.data).as_bytes()
568 );
569 lengths[index] += extent.length;
570 send(publish_content_server_frame::Body::Checkpoint(
574 checkpoint.clone(),
575 ))
576 .await
577 .expect("checkpoint");
578 }
579 Some(publish_content_client_frame::Body::Finish(finish)) => {
580 assert_eq!(finish.checkpoint, Some(checkpoint.clone()));
581 assert_eq!(lengths, [open.packs[0].length, open.packs[1].length]);
582 let mut inventory = Vec::new();
583 for extent in &open.packs {
584 extent
585 .encode_length_delimited(&mut inventory)
586 .expect("inventory");
587 }
588 let mut receipt = PublicationReceipt {
589 client_operation_id: opening.client_operation_id.clone(),
590 destination: open.destination.clone(),
591 thread: open.thread.clone(),
592 revision: open.revision.clone(),
593 sharing_policy_version: if open.sharing_policy_version.is_empty() {
594 vec![3; 32]
595 } else {
596 open.sharing_policy_version.clone()
597 },
598 accepted_inventory: Some(ObjectAddress {
599 algorithm: "blake3".into(),
600 digest: typed_digest("thread-source-inventory-v1", &inventory)
601 .to_vec(),
602 }),
603 outcome: Some(publication_receipt::Outcome::Accepted(
604 Applied::default(),
605 )),
606 import_authority: None,
607 };
608 if wrong_receipt {
609 receipt.thread = None;
610 }
611 let _ =
612 send(publish_content_server_frame::Body::Receipt(receipt)).await;
613 break;
614 }
615 _ => panic!("unexpected frame"),
616 }
617 }
618 });
619 Ok((Writer(Some(tx)), Reader(rx)))
620 }
621 }
622
623 pub(super) fn fixture(
624 wrong_receipt: bool,
625 ) -> (
626 Remote<Peer>,
627 PublishContentClientFrame,
628 [std::io::Cursor<Vec<u8>>; 2],
629 ) {
630 let endpoint = EndpointRef {
631 public_key: vec![8; 32],
632 kind: EndpointKind::Weft as i32,
633 };
634 let remote = Remote {
635 api: api::v2::client::Client::new(
636 Peer {
637 wrong_receipt,
638 original_count: 1,
639 },
640 ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
641 ),
642 description: DescribeEndpointResponse {
643 endpoint: Some(endpoint.clone()),
644 ..Default::default()
645 },
646 };
647 let artifacts = [vec![1; 128 * 1024], vec![2; 16 * 1024]];
648 let packs = artifacts
649 .iter()
650 .zip([
651 pack_extent::Kind::NativePack,
652 pack_extent::Kind::NativeIndex,
653 ])
654 .map(|(data, kind)| {
655 let address = ObjectAddress {
656 algorithm: "blake3".into(),
657 digest: blake3::hash(data).as_bytes().to_vec(),
658 };
659 PackExtent {
660 pack: Some(address.clone()),
661 kind: kind as i32,
662 offset: 0,
663 length: data.len() as u64,
664 extent_digest: Some(address),
665 }
666 })
667 .collect();
668 let open = PublishContentClientFrame {
669 client_operation_id: "op-test".into(),
670 body: Some(publish_content_client_frame::Body::Open(
671 PublishContentOpen {
672 destination: Some(endpoint),
673 thread: Some(ThreadRef::default()),
674 revision: Some(RevisionRef::default()),
675 packs,
676 ..Default::default()
677 },
678 )),
679 };
680 (remote, open, artifacts.map(std::io::Cursor::new))
681 }
682 fn originals() -> PublicationOriginals {
685 let record = SignedRecord {
686 format: "transport-fixture".into(),
687 canonical_record: vec![1],
688 signatures: vec![RecordSignature {
689 public_key: vec![2; 32],
690 signature: vec![3; 64],
691 }],
692 };
693 PublicationOriginals {
694 geneses: vec![ThreadGenesisRecord {
695 boundary_acceptances: Vec::new(),
696 genesis: Some(record.clone()),
697 ..Default::default()
698 }],
699 operations: vec![ReplicationOperations {
700 boundary_acceptances: Vec::new(),
701 operations: vec![record],
702 authority_admissions: vec![],
703 import_authority: None,
704 }],
705 }
706 }
707
708 #[test]
709 fn publication_originals_require_bounded_complete_metadata() {
710 let valid = originals();
711 assert!(valid.validate_bounds().is_ok());
712 let mut missing = valid.clone();
713 missing.operations.clear();
714 assert!(matches!(
715 missing.validate_bounds(),
716 Err(Error::Invalid(
717 "bounded original genesis and source operations required"
718 ))
719 ));
720 let mut unsigned = valid.clone();
721 unsigned.operations[0].operations[0].signatures.clear();
722 assert!(matches!(
723 unsigned.validate_bounds(),
724 Err(Error::Invalid(
725 "original signatures and canonical records required"
726 ))
727 ));
728 let mut oversized = valid.clone();
729 oversized.operations[0].operations[0].canonical_record = vec![1; 256 * 1024];
730 assert!(matches!(
731 oversized.validate_bounds(),
732 Err(Error::Invalid(
733 "publication original metadata budget exceeded"
734 ))
735 ));
736 let mut too_many = valid;
737 too_many.geneses = vec![too_many.geneses[0].clone(); 129];
738 assert!(matches!(
739 too_many.validate_bounds(),
740 Err(Error::Invalid(
741 "bounded original genesis and source operations required"
742 ))
743 ));
744 }
745
746 #[tokio::test]
747 async fn publication_drains_checkpoints_while_uploading_under_backpressure() {
748 let (remote, open, artifacts) = fixture(false);
749 let receipt = tokio::time::timeout(
750 std::time::Duration::from_secs(2),
751 remote.publish_content(&open, &originals(), artifacts),
752 )
753 .await
754 .expect("must not deadlock")
755 .expect("publication");
756 assert_eq!(receipt.client_operation_id, "op-test");
757 }
758 #[tokio::test]
759 async fn publication_rebatches_after_ready_before_upload() {
760 let (mut remote, open, artifacts) = fixture(false);
761 remote.api = api::v2::client::Client::new(
762 Peer {
763 wrong_receipt: false,
764 original_count: 140,
765 },
766 ["/heddle.api.v1alpha2.SyncService/PublishContent".into()],
767 );
768 let mut originals = originals();
769 let mut record = originals.operations[0].operations[0].clone();
770 record.canonical_record = vec![1; 300];
771 originals.operations[0].operations = vec![record; 140];
772 let receipt = tokio::time::timeout(
773 std::time::Duration::from_secs(5),
774 remote.publish_content(&open, &originals, artifacts),
775 )
776 .await
777 .expect("must not deadlock")
778 .expect("negotiated bounded publication");
779 assert_eq!(receipt.client_operation_id, "op-test");
780 }
781
782 #[tokio::test]
783 async fn publication_refuses_originals_carrying_hybrid_import_authority() {
784 let (remote, open, artifacts) = fixture(false);
785 let mut originals = originals();
786 originals.operations[0].import_authority = Some(ImportPublicProofBundleV1::default());
787 let error = remote
788 .publish_content(&open, &originals, artifacts)
789 .await
790 .expect_err("HYBRID originals must not be relayed with the bundle dropped");
791 assert!(
792 matches!(error, Error::Invalid(message) if message.contains("api#307")),
793 "{error}"
794 );
795 }
796
797 #[tokio::test]
798 async fn publication_refuses_a_receipt_for_another_scope() {
799 let (remote, open, artifacts) = fixture(true);
800 let error = remote
801 .publish_content(&open, &originals(), artifacts)
802 .await
803 .expect_err("mismatched scope");
804 assert!(matches!(
805 error,
806 Error::Invalid("receipt differs from requested publication")
807 ));
808 }
809 #[test]
810 fn publication_policy_is_optional_cas_and_receipt_reports_actual_frontier() {
811 let (_, mut opening, _) = fixture(false);
812 let inventory = [7; 32];
813 let Some(publish_content_client_frame::Body::Open(open)) = &opening.body else {
814 panic!("opening")
815 };
816 let receipt = PublicationReceipt {
817 client_operation_id: opening.client_operation_id.clone(),
818 destination: open.destination.clone(),
819 thread: open.thread.clone(),
820 revision: open.revision.clone(),
821 sharing_policy_version: vec![3; 32],
822 accepted_inventory: Some(ObjectAddress {
823 algorithm: "blake3".into(),
824 digest: inventory.to_vec(),
825 }),
826 outcome: Some(publication_receipt::Outcome::Accepted(Applied::default())),
827 import_authority: None,
828 };
829 validate_receipt(receipt.clone(), &opening, &inventory)
830 .expect("one-shot upload returns actual policy without CAS");
831 let mut hybrid = receipt.clone();
832 hybrid.import_authority = Some(ImportPublicProofBundleV1::default());
833 assert!(
834 matches!(
835 validate_receipt(hybrid, &opening, &inventory),
836 Err(Error::Invalid(message)) if message.contains("api#307")
837 ),
838 "a receipt carrying HYBRID import authority must be refused, never ignored"
839 );
840 let mut absent = receipt.clone();
841 absent.sharing_policy_version.clear();
842 assert!(
843 validate_receipt(absent, &opening, &inventory).is_err(),
844 "receipt must report actual policy frontier"
845 );
846 let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
847 panic!("opening")
848 };
849 open.sharing_policy_version = vec![4; 32];
850 assert!(
851 validate_receipt(receipt.clone(), &opening, &inventory).is_err(),
852 "explicit policy CAS cannot silently accept another version"
853 );
854 let Some(publish_content_client_frame::Body::Open(open)) = &mut opening.body else {
855 panic!("opening")
856 };
857 open.sharing_policy_version = vec![3; 32];
858 validate_receipt(receipt, &opening, &inventory).expect("matching explicit policy version");
859 }
860
861 #[cfg(feature = "replication")]
862 #[test]
863 fn publication_digests_match_the_native_typed_hash_format() {
864 assert_eq!(
865 typed_digest("thread-source-inventory-v1", b"bytes"),
866 *heddle_object_model::object::ContentHash::compute_typed(
867 "thread-source-inventory-v1",
868 b"bytes"
869 )
870 .as_bytes()
871 );
872 }
873
874 #[cfg(feature = "source-transfer")]
875 #[tokio::test]
876 async fn thread_publication_prepares_only_selected_source_and_binds_its_revision() {
877 use objects::{
878 object::{Attribution, Blob, Principal, State, Tree, TreeEntry},
879 store::{FsStore, ObjectStore},
880 };
881 let root = tempfile::tempdir().expect("local source scratch");
882 let store = FsStore::new(root.path().join("objects"));
883 store.init().expect("store");
884 let blob = Blob::new(vec![1; 128 * 1024]);
885 store.put_blob(&blob).expect("source");
886 let tree = Tree::from_entries(vec![
887 TreeEntry::file("source.rs", blob.hash(), false).expect("entry"),
888 ]);
889 store.put_tree(&tree).expect("tree");
890 let state = State::new_snapshot(
891 tree.hash(),
892 vec![],
893 Attribution::human(Principal::new("user", "user@example.test")),
894 );
895 let prepared = SourcePack::prepare(
896 &store,
897 &state,
898 root.path(),
899 SourceBudget {
900 max_objects: 16,
901 max_decoded_bytes: 256 * 1024,
902 },
903 )
904 .expect("prepare exact source");
905 let (remote, _, _) = fixture(false);
906 let thread = ThreadRef {
907 spool: Some(SpoolRef {
908 id: "spool-test".into(),
909 }),
910 id: Some(ThreadId { value: vec![1; 32] }),
911 };
912 let receipt = remote
913 .thread(thread.clone())
914 .publish_source(
915 &prepared,
916 &originals(),
917 PublicationOptions {
918 client_operation_id: "source-upload".into(),
919 source: EndpointRef {
920 public_key: vec![2; 32],
921 kind: EndpointKind::Device as i32,
922 },
923 sharing_policy_version: vec![],
924 checkpoint: None,
925 },
926 )
927 .await
928 .expect("one Thread-bound publication");
929 assert_eq!(receipt.thread, Some(thread.clone()));
930 assert_eq!(
931 receipt.revision,
932 Some(RevisionRef {
933 spool: thread.spool,
934 revision: Some(revision_ref::Revision::State(
935 api::heddle::api::common::StateId {
936 value: state.id().as_bytes().to_vec()
937 }
938 ))
939 })
940 );
941 }
942}