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