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 crate::RecordIndexError;
24use crate::StreamRecordIndex;
25use crate::StreamRecordRange;
26
27impl StreamStateMachine {
28    pub fn head(&self, stream_id: &BucketStreamId) -> Option<&StreamMetadata> {
29        self.stream_metadata(stream_id)
30    }
31
32    pub fn record_range(
33        &self,
34        stream_id: &BucketStreamId,
35    ) -> Result<Option<StreamRecordRange>, RecordIndexError> {
36        self.stream_slot(stream_id)
37            .and_then(|slot| slot.record_index.as_ref())
38            .map(StreamRecordIndex::range)
39            .transpose()
40    }
41
42    pub fn offset_for_record(
43        &self,
44        stream_id: &BucketStreamId,
45        record: u64,
46    ) -> Result<Option<u64>, RecordIndexError> {
47        let Some(slot) = self.stream_slot(stream_id) else {
48            return Ok(None);
49        };
50        slot.record_index
51            .as_ref()
52            .map(|index| index.offset_for(record, slot.metadata.tail_offset))
53            .transpose()
54    }
55
56    pub fn record_range_for_append(
57        &self,
58        stream_id: &BucketStreamId,
59        start_offset: u64,
60        next_offset: u64,
61        producer: Option<&ProducerRequest>,
62    ) -> Result<Option<StreamRecordRange>, RecordIndexError> {
63        let Some(slot) = self.stream_slot(stream_id) else {
64            return Ok(None);
65        };
66        if let Some(producer) = producer
67            && let Some(record) = slot.producers.get(&producer.producer_id).and_then(|state| {
68                state.last_items.iter().find(|item| {
69                    item.start_offset == start_offset && item.next_offset == next_offset
70                })
71            })
72            && let (Some(first_record), Some(next_record)) =
73                (record.record_start, record.record_next)
74        {
75            return Ok(Some(StreamRecordRange {
76                first_record,
77                next_record,
78            }));
79        }
80        slot.record_index
81            .as_ref()
82            .map(|index| {
83                Ok(StreamRecordRange {
84                    first_record: index
85                        .record_for_offset(start_offset, slot.metadata.tail_offset)?,
86                    next_record: index.record_for_offset(next_offset, slot.metadata.tail_offset)?,
87                })
88            })
89            .transpose()
90    }
91
92    pub fn stream_attrs(&self, stream_id: &BucketStreamId) -> Option<&StreamAttrs> {
93        self.stream_slot(stream_id)
94            .and_then(|slot| slot.attrs.as_ref())
95    }
96
97    pub fn head_at(&mut self, stream_id: &BucketStreamId, now_ms: u64) -> Option<&StreamMetadata> {
98        self.expire_stream_if_due(stream_id, now_ms);
99        self.stream_metadata(stream_id)
100    }
101
102    pub fn access_requires_write(
103        &self,
104        stream_id: &BucketStreamId,
105        now_ms: u64,
106        renew_ttl: bool,
107    ) -> Result<bool, StreamResponse> {
108        self.validate_stream_scope(stream_id)?;
109        let Some(stream) = self.stream_metadata(stream_id) else {
110            return Err(StreamResponse::error(
111                StreamErrorCode::StreamNotFound,
112                format!("stream '{stream_id}' does not exist"),
113            ));
114        };
115        if stream_is_expired(stream, now_ms) {
116            return Ok(true);
117        }
118        Ok(renew_ttl
119            && stream.stream_ttl_seconds.is_some()
120            && stream.last_ttl_touch_at_ms != now_ms)
121    }
122
123    pub fn hot_start_offset(&self, stream_id: &BucketStreamId) -> u64 {
124        let Some(slot) = self.stream_slot(stream_id) else {
125            return 0;
126        };
127        slot.hot_buffer
128            .first_start_offset()
129            .unwrap_or(slot.metadata.tail_offset)
130    }
131
132    pub fn cold_chunks(&self, stream_id: &BucketStreamId) -> &[ColdChunkRef] {
133        self.stream_slot(stream_id)
134            .map(|slot| slot.cold.cold_chunks())
135            .unwrap_or(&[])
136    }
137
138    pub fn external_segments(&self, stream_id: &BucketStreamId) -> &[ObjectPayloadRef] {
139        self.stream_slot(stream_id)
140            .map(|slot| slot.cold.external_segments())
141            .unwrap_or(&[])
142    }
143
144    pub fn hot_segments(&self, stream_id: &BucketStreamId) -> Vec<HotPayloadSegment> {
145        self.stream_slot(stream_id)
146            .map(|slot| slot.hot_buffer.hot_segments())
147            .unwrap_or_default()
148    }
149
150    pub fn hot_payload_len(&self, stream_id: &BucketStreamId) -> Result<u64, StreamResponse> {
151        let Some(slot) = self.stream_slot(stream_id) else {
152            return Err(StreamResponse::error(
153                StreamErrorCode::StreamNotFound,
154                format!("stream '{stream_id}' does not exist"),
155            ));
156        };
157        Ok(u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64"))
158    }
159
160    pub fn total_hot_payload_bytes(&self) -> u64 {
161        self.registry
162            .slots()
163            .map(|slot| u64::try_from(slot.hot_buffer.len()).expect("payload len fits u64"))
164            .sum()
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 cursor = offset;
269        for (hot_start, hot_segment) in &hot_segments {
270            if cursor >= cold_index_end {
271                break;
272            }
273            let Some(hot_end) = read_segment_end(*hot_start, hot_segment) else {
274                continue;
275            };
276            if hot_end <= cursor {
277                continue;
278            }
279            let gap_end = (*hot_start).min(cold_index_end);
280            push_cold_index_segments(
281                &mut segments,
282                stream_id,
283                slot.cold.cold_generation(),
284                cursor,
285                gap_end,
286            );
287            cursor = cursor.max(hot_end);
288        }
289        push_cold_index_segments(
290            &mut segments,
291            stream_id,
292            slot.cold.cold_generation(),
293            cursor,
294            cold_index_end,
295        );
296        for chunk in self.cold_chunks(stream_id) {
297            let start = offset.max(chunk.start_offset);
298            let end = next_offset.min(chunk.end_offset);
299            if start < end {
300                segments.push((
301                    start,
302                    StreamReadSegment::Object(StreamReadObjectSegment {
303                        object: ObjectPayloadRef::from(chunk),
304                        read_start_offset: start,
305                        len: usize::try_from(end - start).expect("object read len fits usize"),
306                    }),
307                ));
308            }
309        }
310        for object in self.external_segments(stream_id) {
311            let start = offset.max(object.start_offset);
312            let end = next_offset.min(object.end_offset);
313            if start < end {
314                segments.push((
315                    start,
316                    StreamReadSegment::Object(StreamReadObjectSegment {
317                        object: object.clone(),
318                        read_start_offset: start,
319                        len: usize::try_from(end - start).expect("object read len fits usize"),
320                    }),
321                ));
322            }
323        }
324        segments.extend(hot_segments);
325        segments.sort_by_key(|(start, _)| *start);
326        if !segments_cover_range(&segments, offset, next_offset) {
327            return Err(StreamResponse::error_with_next_offset(
328                StreamErrorCode::InvalidColdFlush,
329                format!("stream '{stream_id}' has missing payload segment metadata"),
330                next_offset,
331            ));
332        }
333        Ok(StreamReadPlan {
334            offset,
335            next_offset,
336            content_type: stream.content_type.clone(),
337            segments: segments.into_iter().map(|(_, segment)| segment).collect(),
338            up_to_date: next_offset == stream.tail_offset,
339            closed: stream.status == StreamStatus::Closed,
340            retained_record_range: None,
341            record_range: None,
342        })
343    }
344
345    pub fn latest_snapshot(
346        &self,
347        stream_id: &BucketStreamId,
348    ) -> Result<Option<StreamVisibleSnapshot>, StreamResponse> {
349        let Some(slot) = self.stream_slot(stream_id) else {
350            return Err(StreamResponse::error(
351                StreamErrorCode::StreamNotFound,
352                format!("stream '{stream_id}' does not exist"),
353            ));
354        };
355        Ok(slot.visible_snapshot.clone())
356    }
357
358    pub fn read_snapshot(
359        &self,
360        stream_id: &BucketStreamId,
361        snapshot_offset: u64,
362    ) -> Result<StreamVisibleSnapshot, StreamResponse> {
363        let snapshot = self.latest_snapshot(stream_id)?;
364        match snapshot {
365            Some(snapshot) if snapshot.offset == snapshot_offset => Ok(snapshot),
366            _ => Err(StreamResponse::error(
367                StreamErrorCode::SnapshotNotFound,
368                format!("snapshot {snapshot_offset} for stream '{stream_id}' does not exist"),
369            )),
370        }
371    }
372
373    pub fn bootstrap_plan(
374        &self,
375        stream_id: &BucketStreamId,
376    ) -> Result<StreamBootstrapPlan, StreamResponse> {
377        let Some(slot) = self.stream_slot(stream_id) else {
378            return Err(StreamResponse::error(
379                StreamErrorCode::StreamNotFound,
380                format!("stream '{stream_id}' does not exist"),
381            ));
382        };
383        let stream = &slot.metadata;
384        let snapshot = slot.visible_snapshot.clone();
385        let retained_offset = snapshot
386            .as_ref()
387            .map(|snapshot| snapshot.offset)
388            .unwrap_or(0);
389        let updates = slot
390            .message_records
391            .iter()
392            .filter(|record| record.start_offset >= retained_offset)
393            .cloned()
394            .collect::<Vec<_>>();
395        Ok(StreamBootstrapPlan {
396            snapshot,
397            updates,
398            next_offset: stream.tail_offset,
399            content_type: stream.content_type.clone(),
400            up_to_date: true,
401            closed: stream.status == StreamStatus::Closed,
402        })
403    }
404}
405
406fn push_cold_index_segments(
407    segments: &mut Vec<(u64, StreamReadSegment)>,
408    _stream_id: &BucketStreamId,
409    generation: u64,
410    start_offset: u64,
411    end_offset: u64,
412) {
413    let mut cursor = start_offset;
414    while cursor < end_offset {
415        let page_id = cursor / COLD_INDEX_PAGE_SPAN_BYTES;
416        let page_end = page_id
417            .saturating_add(1)
418            .saturating_mul(COLD_INDEX_PAGE_SPAN_BYTES);
419        let segment_end = end_offset.min(page_end);
420        segments.push((
421            cursor,
422            StreamReadSegment::ColdIndex(StreamReadColdIndexSegment {
423                generation,
424                page_id,
425                read_start_offset: cursor,
426                len: usize::try_from(segment_end - cursor).expect("cold index read len fits usize"),
427            }),
428        ));
429        cursor = segment_end;
430    }
431}
432
433fn segments_cover_range(
434    segments: &[(u64, StreamReadSegment)],
435    offset: u64,
436    next_offset: u64,
437) -> bool {
438    if next_offset < offset {
439        return false;
440    }
441    let mut expected_start = offset;
442    for (segment_start, segment) in segments {
443        let Some(segment_end) = read_segment_end(*segment_start, segment) else {
444            return false;
445        };
446        if segment_end <= expected_start {
447            continue;
448        }
449        if *segment_start > expected_start {
450            return false;
451        }
452        expected_start = segment_end;
453        if expected_start >= next_offset {
454            return true;
455        }
456    }
457    expected_start == next_offset
458}
459
460fn read_segment_end(segment_start: u64, segment: &StreamReadSegment) -> Option<u64> {
461    match segment {
462        StreamReadSegment::Object(object) => {
463            if object.len == 0
464                || object.read_start_offset != segment_start
465                || object.read_start_offset < object.object.start_offset
466            {
467                return None;
468            }
469            let len = u64::try_from(object.len).ok()?;
470            let segment_end = object.read_start_offset.checked_add(len)?;
471            if segment_end > object.object.end_offset {
472                return None;
473            }
474            Some(segment_end)
475        }
476        StreamReadSegment::ColdIndex(index) => {
477            if index.len == 0 || index.read_start_offset != segment_start {
478                return None;
479            }
480            let len = u64::try_from(index.len).ok()?;
481            segment_start.checked_add(len)
482        }
483        StreamReadSegment::Hot(payload) => {
484            if payload.is_empty() {
485                return None;
486            }
487            let len = u64::try_from(payload.len()).ok()?;
488            segment_start.checked_add(len)
489        }
490    }
491}