Skip to main content

ursula_stream/state_machine/
query.rs

1//! Read and query paths: heads, attrs, hot/cold accessors, read plans, snapshots, bootstrap.
2
3use super::BucketStreamId;
4use super::COLD_INDEX_PAGE_SPAN_BYTES;
5use super::ColdChunkRef;
6use super::HotPayloadSegment;
7use super::ObjectPayloadRef;
8use super::StreamAttrs;
9use super::StreamBootstrapPlan;
10use super::StreamErrorCode;
11use super::StreamMetadata;
12use super::StreamRead;
13use super::StreamReadColdIndexSegment;
14use super::StreamReadObjectSegment;
15use super::StreamReadPlan;
16use super::StreamReadSegment;
17use super::StreamResponse;
18use super::StreamStateMachine;
19use super::StreamStatus;
20use super::StreamVisibleSnapshot;
21use super::is_soft_deleted;
22use super::stream_is_expired;
23
24impl StreamStateMachine {
25    pub fn head(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
26        self.stream_metadata(stream_id)
27    }
28
29    pub fn stream_attrs(&self, stream_id: &BucketStreamId) -> Option<&StreamAttrs> {
30        self.stream_slot(stream_id)
31            .and_then(|slot| slot.attrs.as_ref())
32    }
33
34    pub fn head_at(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> Option<&StreamMetadata> {
35        self.expire_stream_if_due(stream_id, now_ms);
36        self.stream_metadata(stream_id)
37    }
38
39    pub fn access_requires_write(
40        &self,
41        stream_id: &BucketStreamId,
42        now_ms: u64,
43        renew_ttl: bool,
44    ) -> Result<bool, StreamResponse> {
45        self.validate_stream_scope(stream_id)?;
46        let Some(stream) = self.stream_metadata(stream_id) else {
47            return Err(StreamResponse::error(
48                StreamErrorCode::StreamNotFound,
49                format!("stream '{stream_id}' does not exist"),
50            ));
51        };
52        if is_soft_deleted(stream) {
53            return Err(StreamResponse::error(
54                StreamErrorCode::StreamGone,
55                format!("stream '{stream_id}' is gone"),
56            ));
57        }
58        if stream_is_expired(stream, now_ms) {
59            return Ok(true);
60        }
61        Ok(renew_ttl
62            && stream.stream_ttl_seconds.is_some()
63            && stream.last_ttl_touch_at_ms != now_ms)
64    }
65
66    pub fn hot_start_offset(&self, stream_id: &BucketStreamId) -> u64 {
67        let Some(slot) = self.stream_slot(stream_id) else {
68            return 0;
69        };
70        slot.hot_buffer
71            .first_start_offset()
72            .unwrap_or(slot.metadata.tail_offset)
73    }
74
75    pub fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
76        self.stream_slot(stream_id)
77            .map(|slot| slot.cold.cold_chunks())
78            .unwrap_or(&[])
79    }
80
81    pub fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
82        self.stream_slot(stream_id)
83            .map(|slot| slot.cold.external_segments())
84            .unwrap_or(&[])
85    }
86
87    pub fn hot_segments(&self, stream_id: &BucketStreamId) -> Vec<HotPayloadSegment> {
88        self.stream_slot(stream_id)
89            .map(|slot| slot.hot_buffer.hot_segments())
90            .unwrap_or_default()
91    }
92
93    pub fn hot_payload_len(&self, stream_id: &BucketStreamId) -> Result<u64, StreamResponse> {
94        let Some(slot) = self.stream_slot(stream_id) else {
95            return Err(StreamResponse::error(
96                StreamErrorCode::StreamNotFound,
97                format!("stream '{stream_id}' does not exist"),
98            ));
99        };
100        if is_soft_deleted(&slot.metadata) {
101            return Err(StreamResponse::error(
102                StreamErrorCode::StreamGone,
103                format!("stream '{stream_id}' is gone"),
104            ));
105        }
106        Ok(u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64"))
107    }
108
109    pub fn total_hot_payload_bytes(&self) -> u64 {
110        self.registry
111            .slots()
112            .map(|slot| u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64"))
113            .sum()
114    }
115
116    pub fn bucket_exists(&self, bucket_id: &str) -> bool {
117        self.buckets.contains(bucket_id)
118    }
119
120    pub fn read(
121        &self,
122        stream_id: &BucketStreamId,
123        offset: u64,
124        max_len: usize,
125    ) -> Result<StreamRead, StreamResponse> {
126        let plan = self.read_plan(stream_id, offset, max_len)?;
127        if plan.segments.iter().any(|segment| {
128            matches!(
129                segment,
130                StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_)
131            )
132        }) {
133            return Err(StreamResponse::error_with_next_offset(
134                StreamErrorCode::InvalidColdFlush,
135                format!("stream '{stream_id}' read requires object payload store"),
136                plan.next_offset,
137            ));
138        }
139        let payload = plan
140            .segments
141            .iter()
142            .flat_map(|segment| match segment {
143                StreamReadSegment::Hot(payload) => payload.as_slice(),
144                StreamReadSegment::ColdIndex(_) | StreamReadSegment::Object(_) => {
145                    unreachable!("object segments checked above")
146                }
147            })
148            .copied()
149            .collect();
150        Ok(StreamRead {
151            offset: plan.offset,
152            next_offset: plan.next_offset,
153            content_type: plan.content_type,
154            payload,
155            up_to_date: plan.up_to_date,
156            closed: plan.closed,
157        })
158    }
159
160    pub fn read_plan(
161        &self,
162        stream_id: &BucketStreamId,
163        offset: u64,
164        max_len: usize,
165    ) -> Result<StreamReadPlan, StreamResponse> {
166        self.read_plan_at(stream_id, offset, max_len, 0)
167    }
168
169    pub fn read_plan_at(
170        &self,
171        stream_id: &BucketStreamId,
172        offset: u64,
173        max_len: usize,
174        now_ms: u64,
175    ) -> Result<StreamReadPlan, StreamResponse> {
176        let Some(slot) = self.stream_slot(stream_id) else {
177            return Err(StreamResponse::error(
178                StreamErrorCode::StreamNotFound,
179                format!("stream '{stream_id}' does not exist"),
180            ));
181        };
182        let stream = &slot.metadata;
183        if is_soft_deleted(stream) {
184            return Err(StreamResponse::error(
185                StreamErrorCode::StreamGone,
186                format!("stream '{stream_id}' is gone"),
187            ));
188        }
189        if stream_is_expired(stream, now_ms) {
190            return Err(StreamResponse::error(
191                StreamErrorCode::StreamNotFound,
192                format!("stream '{stream_id}' does not exist"),
193            ));
194        }
195        if offset > stream.tail_offset {
196            return Err(StreamResponse::error_with_next_offset(
197                StreamErrorCode::OffsetOutOfRange,
198                format!(
199                    "offset {offset} is beyond stream '{}' tail {}",
200                    stream_id, stream.tail_offset
201                ),
202                stream.tail_offset,
203            ));
204        }
205        let retained_offset = self.earliest_retained_offset(stream_id);
206        if offset < retained_offset {
207            return Err(StreamResponse::error_with_next_offset(
208                StreamErrorCode::StreamGone,
209                format!(
210                    "offset {offset} is older than stream '{}' retained offset {retained_offset}",
211                    stream_id
212                ),
213                retained_offset,
214            ));
215        }
216
217        let max_len_u64 = u64::try_from(max_len).unwrap_or(u64::MAX);
218        let next_offset = stream.tail_offset.min(offset.saturating_add(max_len_u64));
219        let mut segments = Vec::<(u64, StreamReadSegment)>::new();
220        let hot_segments = slot.hot_buffer.read_segments(offset, next_offset);
221        let cold_frontier = self.cold_frontier_offset(stream_id, retained_offset);
222        let cold_index_end = next_offset.min(cold_frontier);
223        let mut cursor = offset;
224        for (hot_start, hot_segment) in &hot_segments {
225            if cursor >= cold_index_end {
226                break;
227            }
228            let Some(hot_end) = read_segment_end(*hot_start, hot_segment) else {
229                continue;
230            };
231            if hot_end <= cursor {
232                continue;
233            }
234            let gap_end = (*hot_start).min(cold_index_end);
235            push_cold_index_segments(
236                &mut segments,
237                stream_id,
238                slot.cold.cold_generation(),
239                cursor,
240                gap_end,
241            );
242            cursor = cursor.max(hot_end);
243        }
244        push_cold_index_segments(
245            &mut segments,
246            stream_id,
247            slot.cold.cold_generation(),
248            cursor,
249            cold_index_end,
250        );
251        for chunk in self.cold_chunks(stream_id) {
252            let start = offset.max(chunk.start_offset);
253            let end = next_offset.min(chunk.end_offset);
254            if start < end {
255                segments.push((
256                    start,
257                    StreamReadSegment::Object(StreamReadObjectSegment {
258                        object: ObjectPayloadRef::from(chunk),
259                        read_start_offset: start,
260                        len: usize::try_from(end - start).expect("object read len fits usize"),
261                    }),
262                ));
263            }
264        }
265        for object in self.external_segments(stream_id) {
266            let start = offset.max(object.start_offset);
267            let end = next_offset.min(object.end_offset);
268            if start < end {
269                segments.push((
270                    start,
271                    StreamReadSegment::Object(StreamReadObjectSegment {
272                        object: object.clone(),
273                        read_start_offset: start,
274                        len: usize::try_from(end - start).expect("object read len fits usize"),
275                    }),
276                ));
277            }
278        }
279        segments.extend(hot_segments);
280        segments.sort_by_key(|(start, _)| *start);
281        if !segments_cover_range(&segments, offset, next_offset) {
282            return Err(StreamResponse::error_with_next_offset(
283                StreamErrorCode::InvalidColdFlush,
284                format!("stream '{stream_id}' has missing payload segment metadata"),
285                next_offset,
286            ));
287        }
288        Ok(StreamReadPlan {
289            offset,
290            next_offset,
291            content_type: stream.content_type.clone(),
292            segments: segments.into_iter().map(|(_, segment)| segment).collect(),
293            up_to_date: next_offset == stream.tail_offset,
294            closed: stream.status == StreamStatus::Closed,
295        })
296    }
297
298    pub fn latest_snapshot(
299        &self,
300        stream_id: &BucketStreamId,
301    ) -> Result<Option<StreamVisibleSnapshot>, StreamResponse> {
302        let Some(slot) = self.stream_slot(stream_id) else {
303            return Err(StreamResponse::error(
304                StreamErrorCode::StreamNotFound,
305                format!("stream '{stream_id}' does not exist"),
306            ));
307        };
308        if is_soft_deleted(&slot.metadata) {
309            return Err(StreamResponse::error(
310                StreamErrorCode::StreamGone,
311                format!("stream '{stream_id}' is gone"),
312            ));
313        }
314        Ok(slot.visible_snapshot.clone())
315    }
316
317    pub fn read_snapshot(
318        &self,
319        stream_id: &BucketStreamId,
320        snapshot_offset: u64,
321    ) -> Result<StreamVisibleSnapshot, StreamResponse> {
322        let snapshot = self.latest_snapshot(stream_id)?;
323        match snapshot {
324            Some(snapshot) if snapshot.offset == snapshot_offset => Ok(snapshot),
325            _ => Err(StreamResponse::error(
326                StreamErrorCode::SnapshotNotFound,
327                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
328            )),
329        }
330    }
331
332    pub fn bootstrap_plan(
333        &self,
334        stream_id: &BucketStreamId,
335    ) -> Result<StreamBootstrapPlan, StreamResponse> {
336        let Some(slot) = self.stream_slot(stream_id) else {
337            return Err(StreamResponse::error(
338                StreamErrorCode::StreamNotFound,
339                format!("stream '{stream_id}' does not exist"),
340            ));
341        };
342        let stream = &slot.metadata;
343        if is_soft_deleted(stream) {
344            return Err(StreamResponse::error(
345                StreamErrorCode::StreamGone,
346                format!("stream '{stream_id}' is gone"),
347            ));
348        }
349        let snapshot = slot.visible_snapshot.clone();
350        let retained_offset = snapshot
351            .as_ref()
352            .map(|snapshot| snapshot.offset)
353            .unwrap_or(0);
354        let updates = slot
355            .message_records
356            .iter()
357            .filter(|record| record.start_offset >= retained_offset)
358            .cloned()
359            .collect::<Vec<_>>();
360        Ok(StreamBootstrapPlan {
361            snapshot,
362            updates,
363            next_offset: stream.tail_offset,
364            content_type: stream.content_type.clone(),
365            up_to_date: true,
366            closed: stream.status == StreamStatus::Closed,
367        })
368    }
369}
370
371fn push_cold_index_segments(
372    segments: &mut Vec<(u64, StreamReadSegment)>,
373    _stream_id: &BucketStreamId,
374    generation: u64,
375    start_offset: u64,
376    end_offset: u64,
377) {
378    let mut cursor = start_offset;
379    while cursor < end_offset {
380        let page_id = cursor / COLD_INDEX_PAGE_SPAN_BYTES;
381        let page_end = page_id
382            .saturating_add(1)
383            .saturating_mul(COLD_INDEX_PAGE_SPAN_BYTES);
384        let segment_end = end_offset.min(page_end);
385        segments.push((
386            cursor,
387            StreamReadSegment::ColdIndex(StreamReadColdIndexSegment {
388                generation,
389                page_id,
390                read_start_offset: cursor,
391                len: usize::try_from(segment_end - cursor).expect("cold index read len fits usize"),
392            }),
393        ));
394        cursor = segment_end;
395    }
396}
397
398fn segments_cover_range(
399    segments: &[(u64, StreamReadSegment)],
400    offset: u64,
401    next_offset: u64,
402) -> bool {
403    if next_offset < offset {
404        return false;
405    }
406    let mut expected_start = offset;
407    for (segment_start, segment) in segments {
408        let Some(segment_end) = read_segment_end(*segment_start, segment) else {
409            return false;
410        };
411        if segment_end <= expected_start {
412            continue;
413        }
414        if *segment_start > expected_start {
415            return false;
416        }
417        expected_start = segment_end;
418        if expected_start >= next_offset {
419            return true;
420        }
421    }
422    expected_start == next_offset
423}
424
425fn read_segment_end(segment_start: u64, segment: &StreamReadSegment) -> Option<u64> {
426    match segment {
427        StreamReadSegment::Object(object) => {
428            if object.len == 0
429                || object.read_start_offset != segment_start
430                || object.read_start_offset < object.object.start_offset
431            {
432                return None;
433            }
434            let len = u64::try_from(object.len).ok()?;
435            let segment_end = object.read_start_offset.checked_add(len)?;
436            if segment_end > object.object.end_offset {
437                return None;
438            }
439            Some(segment_end)
440        }
441        StreamReadSegment::ColdIndex(index) => {
442            if index.len == 0 || index.read_start_offset != segment_start {
443                return None;
444            }
445            let len = u64::try_from(index.len).ok()?;
446            segment_start.checked_add(len)
447        }
448        StreamReadSegment::Hot(payload) => {
449            if payload.is_empty() {
450                return None;
451            }
452            let len = u64::try_from(payload.len()).ok()?;
453            segment_start.checked_add(len)
454        }
455    }
456}