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