1use 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}