heddle_thread_api/
content.rs1use api::v2::client::RpcTransport;
5
6pub use crate::contract::blob_read::Source as BlobSource;
7use crate::{
8 Remote,
9 contract::*,
10 observation::{self, Error},
11 rpc, transport,
12};
13
14pub struct Blob {
15 pub source: BlobSource,
16 pub object_hash: Vec<u8>,
17 pub bytes: Vec<u8>,
18}
19
20impl<T: RpcTransport<Error = transport::Error>> Remote<T> {
21 pub async fn read_blobs(
25 &self,
26 thread: ThreadRef,
27 revision: RevisionRef,
28 sources: Vec<BlobSource>,
29 ) -> Result<Vec<Blob>, Error> {
30 if thread.spool != revision.spool
31 || thread.id.as_ref().is_none_or(|id| id.value.len() != 32)
32 {
33 return Err(Error::Invalid(
34 "content Thread differs from exact revision scope",
35 ));
36 }
37 let budget = observation::budget(&self.description)?;
38 if sources.is_empty() || sources.len() > budget.max_items as usize {
39 return Err(Error::Invalid("invalid selection count"));
40 }
41 for source in &sources {
42 match source {
43 BlobSource::Path(path) if !path.is_empty() => {}
44 BlobSource::ObjectHash(hash) if hash.len() == 32 => {}
45 _ => return Err(Error::Invalid("invalid blob source")),
46 }
47 }
48 crate::reopen::retry(|| {
49 self.read_blobs_once(thread.clone(), revision.clone(), sources.clone(), budget)
50 })
51 .await
52 }
53
54 async fn read_blobs_once(
55 &self,
56 thread: ThreadRef,
57 revision: RevisionRef,
58 sources: Vec<BlobSource>,
59 budget: ReadBudget,
60 ) -> Result<Vec<Blob>, Error> {
61 let selections = sources
62 .iter()
63 .enumerate()
64 .map(|(i, source)| ContentRead {
65 selection_id: i.to_string(),
66 selection: Some(content_read::Selection::Blob(BlobRead {
67 source: Some(source.clone()),
68 offset: 0,
69 length: 0,
70 })),
71 })
72 .collect();
73 let mut messages = self
74 .api
75 .observe::<rpc::ContentServiceReadContent>(&ReadContentRequest {
76 thread: Some(thread),
77 revision: Some(revision.clone()),
78 selections,
79 budget: Some(budget),
80 })
81 .await?;
82 let mut blobs: Vec<_> = sources
83 .into_iter()
84 .map(|source| Blob {
85 source,
86 object_hash: vec![],
87 bytes: vec![],
88 })
89 .collect();
90 let mut range_done = vec![false; blobs.len()];
91 let mut complete = vec![false; blobs.len()];
92 let mut totals = vec![None; blobs.len()];
93 let mut accepted = None;
94 let mut received_bytes = 0_u64;
95 let mut received_items = 0_u32;
96 while let Some(event) = messages.next().await? {
98 let echo = matches!(
99 event.payload,
100 Some(content_event::Payload::AcceptedBudget(_))
101 );
102 if let Some(content_event::Payload::AcceptedBudget(effective)) = &event.payload {
103 if accepted.is_some() || !event.selection_id.is_empty() {
104 return Err(Error::Invalid(
105 "duplicate or selection-scoped content budget echo",
106 ));
107 }
108 if event
109 .revision
110 .as_ref()
111 .is_some_and(|echo_revision| echo_revision != &revision)
112 {
113 return Err(Error::Invalid("content revision mismatch"));
114 }
115 api::v2::validate_accepted_read_budget(&budget, Some(effective))
116 .map_err(|_| Error::Invalid("invalid accepted content budget"))?;
117 accepted = Some(*effective);
118 }
119 let limits = accepted.ok_or(Error::Invalid("missing initial content budget echo"))?;
120 let size = prost::Message::encoded_len(&event) as u64 + 5;
123 if size > u64::from(limits.max_frame_bytes)
124 || size > limits.max_snapshot_bytes.saturating_sub(received_bytes)
125 || received_items >= limits.max_items
126 {
127 return Err(Error::Invalid("content budget exceeded"));
128 }
129 received_bytes += size;
130 received_items += 1;
131 if echo {
132 continue;
133 }
134 if event.revision.as_ref() != Some(&revision) {
135 return Err(Error::Invalid("content revision mismatch"));
136 }
137 let index = event
138 .selection_id
139 .parse::<usize>()
140 .map_err(|_| Error::Invalid("unknown content selection"))?;
141 if event.selection_id != index.to_string() || index >= blobs.len() || complete[index] {
142 return Err(Error::Invalid("unknown or completed content selection"));
143 }
144 let blob = &mut blobs[index];
145 match event
146 .payload
147 .ok_or(Error::Invalid("missing content payload"))?
148 {
149 content_event::Payload::Blob(chunk) => {
150 if range_done[index]
151 || chunk.offset != blob.bytes.len() as u64
152 || chunk.total_size > limits.max_snapshot_bytes
153 || chunk.data.len() as u64 > chunk.total_size.saturating_sub(chunk.offset)
154 || chunk.object_hash.len() != 32
155 || totals[index].is_some_and(|total| total != chunk.total_size)
156 || (!blob.object_hash.is_empty() && blob.object_hash != chunk.object_hash)
157 || matches!(&blob.source, BlobSource::ObjectHash(hash) if *hash != chunk.object_hash)
158 {
159 return Err(Error::Invalid("inconsistent blob range or identity"));
160 }
161 totals[index] = Some(chunk.total_size);
162 blob.object_hash = chunk.object_hash;
163 blob.bytes.extend(chunk.data);
164 if chunk.range_complete && blob.bytes.len() as u64 != chunk.total_size {
165 return Err(Error::Invalid("truncated complete blob"));
166 }
167 range_done[index] = chunk.range_complete;
168 }
169 content_event::Payload::SelectionComplete(status) => {
170 if !range_done[index]
171 || status.coverage != Coverage::Complete as i32
172 || status.computed_for.as_ref().is_some_and(|r| r != &revision)
173 {
174 return Err(Error::Invalid("incomplete blob selection"));
175 }
176 complete[index] = true;
177 }
178 _ => return Err(Error::Invalid("unexpected content payload")),
179 }
180 }
181 if accepted.is_none() {
182 return Err(Error::Invalid("missing initial content budget echo"));
183 }
184 if complete.iter().all(|done| *done) {
185 Ok(blobs)
186 } else {
187 Err(Error::Interrupted)
188 }
189 }
190}
191
192#[cfg(feature = "replication")]
195pub fn structured_conflicts(
196 bytes: &[u8],
197) -> Result<heddle_object_model::object::StructuredConflict, transport::Error> {
198 heddle_object_model::object::StructuredConflict::decode(bytes)
199 .map_err(|error| transport::Error::Io(error.to_string()))
200}