ursula_stream/state_machine/
query.rs1use 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}