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