Skip to main content

heddle_thread_api/
content.rs

1// SPDX-License-Identifier: Apache-2.0
2//! Bounded full-blob reads at one exact revision. Large/ranged transfers use the
3//! same generated ContentServiceReadContent stream and a caller-owned sink.
4use 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    /// One request for any mixture of paths and missing object hashes. The
22    /// The selected Thread and revision pin content and authorization without a
23    /// mutable Thread-tip lookup.
24    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        // Validate through FIN so selection completion cannot hide a trailing echo.
97        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            // Hosted stream framing is a one-byte kind and four-byte body length.
121            // The echo and selection completions consume the same budget as blobs.
122            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/// Decode native conflict attachment bytes without inventing lifecycle evidence.
193/// Region geometry is immutable; absent retained resolution evidence stays unspecified.
194#[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}