Skip to main content

ursula_stream/state_machine/
append.rs

1//! Append paths (inline/external/batch) and idempotent producer bookkeeping.
2
3use super::AppendExternalInput;
4use super::AppendStreamInput;
5use super::BucketStreamId;
6use super::ObjectPayloadRef;
7use super::ProducerAppendRecord;
8use super::ProducerRequest;
9use super::ProducerState;
10use super::StreamBatchAppend;
11use super::StreamBatchAppendItem;
12use super::StreamErrorCode;
13use super::StreamErrorContext;
14use super::StreamMessageRecord;
15use super::StreamMetadata;
16use super::StreamResponse;
17use super::StreamStateMachine;
18use super::StreamStatus;
19use super::is_soft_deleted;
20use super::renew_stream_ttl;
21use super::validate_external_payload_ref;
22use super::validate_producer_request;
23
24impl StreamStateMachine {
25    pub fn append_borrowed(&mut self, input: AppendStreamInput<'_>) -> StreamResponse {
26        let AppendStreamInput {
27            stream_id,
28            content_type,
29            payload,
30            close_after,
31            stream_seq,
32            producer,
33            now_ms,
34        } = input;
35        if let Err(response) = self.validate_stream_scope(&stream_id) {
36            return response;
37        }
38        if let Err(response) = validate_producer_request(producer.as_ref()) {
39            return response;
40        }
41
42        let Some(_) = self.stream_metadata(&stream_id) else {
43            return StreamResponse::error(
44                StreamErrorCode::StreamNotFound,
45                format!("stream '{stream_id}' does not exist"),
46            );
47        };
48        if self.expire_stream_if_due(&stream_id, now_ms) {
49            return StreamResponse::error(
50                StreamErrorCode::StreamNotFound,
51                format!("stream '{stream_id}' does not exist"),
52            );
53        }
54        if self
55            .stream_metadata(&stream_id)
56            .is_some_and(is_soft_deleted)
57        {
58            return StreamResponse::error(
59                StreamErrorCode::StreamGone,
60                format!("stream '{stream_id}' is gone"),
61            );
62        }
63        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
64            Ok(decision) => decision,
65            Err(response) => return response,
66        };
67        if let ProducerDecision::Duplicate {
68            offset,
69            next_offset,
70            closed,
71            producer,
72            ..
73        } = producer_decision
74        {
75            if payload.is_empty() {
76                return StreamResponse::Closed {
77                    next_offset,
78                    deduplicated: true,
79                    producer: Some(producer),
80                };
81            }
82            return StreamResponse::Appended {
83                offset,
84                next_offset,
85                closed,
86                deduplicated: true,
87                producer: Some(producer),
88            };
89        }
90
91        let Some(stream) = self.stream_metadata_mut(&stream_id) else {
92            unreachable!("stream existence checked before producer evaluation");
93        };
94
95        if stream.status == StreamStatus::Closed {
96            if close_after && payload.is_empty() {
97                return StreamResponse::Closed {
98                    next_offset: stream.tail_offset,
99                    deduplicated: false,
100                    producer: None,
101                };
102            }
103            return StreamResponse::error_with_next_offset_and_context(
104                StreamErrorCode::StreamClosed,
105                format!("stream '{stream_id}' is closed"),
106                stream.tail_offset,
107                vec![StreamErrorContext::StreamClosed],
108            );
109        }
110
111        if payload.is_empty() && !close_after {
112            return StreamResponse::error(
113                StreamErrorCode::EmptyAppend,
114                "append payload must be non-empty unless closing the stream",
115            );
116        }
117
118        if !payload.is_empty() {
119            let Some(content_type) = content_type else {
120                return StreamResponse::error(
121                    StreamErrorCode::MissingContentType,
122                    "append with a body must include content type",
123                );
124            };
125            if content_type != stream.content_type {
126                return StreamResponse::error_with_next_offset(
127                    StreamErrorCode::ContentTypeMismatch,
128                    format!(
129                        "append content type '{content_type}' does not match stream content type '{}'",
130                        stream.content_type
131                    ),
132                    stream.tail_offset,
133                );
134            }
135        }
136
137        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
138            return response;
139        }
140
141        let offset = stream.tail_offset;
142        let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
143        stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
144        if let Some(seq) = stream_seq {
145            stream.last_stream_seq = Some(seq);
146        }
147        renew_stream_ttl(stream, now_ms);
148        if close_after {
149            stream.status = StreamStatus::Closed;
150        }
151        let closed = stream.status == StreamStatus::Closed;
152        let next_offset = stream.tail_offset;
153        self.refresh_ttl_entry(&stream_id);
154        let producer_ack = producer.clone();
155        if let Some(producer) = producer {
156            self.record_producer_success(
157                stream_id.clone(),
158                producer,
159                ProducerAppendRecord {
160                    start_offset: offset,
161                    next_offset,
162                    closed,
163                },
164                vec![ProducerAppendRecord {
165                    start_offset: offset,
166                    next_offset,
167                    closed,
168                }],
169            );
170        }
171
172        if payload.is_empty() {
173            StreamResponse::Closed {
174                next_offset,
175                deduplicated: false,
176                producer: producer_ack,
177            }
178        } else {
179            let slot = self
180                .stream_slot_mut(&stream_id)
181                .expect("stream existence checked before append mutation");
182            slot.hot_buffer.push(offset, next_offset, payload);
183            slot.integrity
184                .append_payload(&stream_id, offset, next_offset, payload);
185            slot.message_records.push(StreamMessageRecord {
186                start_offset: offset,
187                end_offset: next_offset,
188            });
189            StreamResponse::Appended {
190                offset,
191                next_offset,
192                closed: close_after,
193                deduplicated: false,
194                producer: producer_ack,
195            }
196        }
197    }
198
199    pub(super) fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
200        let AppendExternalInput {
201            stream_id,
202            content_type,
203            payload,
204            close_after,
205            stream_seq,
206            producer,
207            now_ms,
208        } = input;
209        if let Err(response) = validate_external_payload_ref(&payload) {
210            return response;
211        }
212        if let Err(response) = self.validate_stream_scope(&stream_id) {
213            return response;
214        }
215        if let Err(response) = validate_producer_request(producer.as_ref()) {
216            return response;
217        }
218        let Some(_) = self.stream_metadata(&stream_id) else {
219            return StreamResponse::error(
220                StreamErrorCode::StreamNotFound,
221                format!("stream '{stream_id}' does not exist"),
222            );
223        };
224        if self.expire_stream_if_due(&stream_id, now_ms) {
225            return StreamResponse::error(
226                StreamErrorCode::StreamNotFound,
227                format!("stream '{stream_id}' does not exist"),
228            );
229        }
230        if self
231            .stream_metadata(&stream_id)
232            .is_some_and(is_soft_deleted)
233        {
234            return StreamResponse::error(
235                StreamErrorCode::StreamGone,
236                format!("stream '{stream_id}' is gone"),
237            );
238        }
239        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
240            Ok(decision) => decision,
241            Err(response) => return response,
242        };
243        if let ProducerDecision::Duplicate {
244            offset,
245            next_offset,
246            closed,
247            producer,
248            ..
249        } = producer_decision
250        {
251            return StreamResponse::Appended {
252                offset,
253                next_offset,
254                closed,
255                deduplicated: true,
256                producer: Some(producer),
257            };
258        }
259
260        let Some(stream) = self.stream_metadata(&stream_id) else {
261            unreachable!("stream existence checked before producer evaluation");
262        };
263        if stream.status == StreamStatus::Closed {
264            return StreamResponse::error_with_next_offset_and_context(
265                StreamErrorCode::StreamClosed,
266                format!("stream '{stream_id}' is closed"),
267                stream.tail_offset,
268                vec![StreamErrorContext::StreamClosed],
269            );
270        }
271        let Some(content_type) = content_type else {
272            return StreamResponse::error(
273                StreamErrorCode::MissingContentType,
274                "append with a body must include content type",
275            );
276        };
277        if content_type != stream.content_type {
278            return StreamResponse::error_with_next_offset(
279                StreamErrorCode::ContentTypeMismatch,
280                format!(
281                    "append content type '{content_type}' does not match stream content type '{}'",
282                    stream.content_type
283                ),
284                stream.tail_offset,
285            );
286        }
287        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
288            return response;
289        }
290        let offset = stream.tail_offset;
291        let next_offset = offset.saturating_add(payload.payload_len);
292        let stream = self
293            .stream_metadata_mut(&stream_id)
294            .expect("stream existence checked before external append mutation");
295        stream.tail_offset = next_offset;
296        if let Some(seq) = stream_seq {
297            stream.last_stream_seq = Some(seq);
298        }
299        renew_stream_ttl(stream, now_ms);
300        if close_after {
301            stream.status = StreamStatus::Closed;
302        }
303        let closed = stream.status == StreamStatus::Closed;
304        self.refresh_ttl_entry(&stream_id);
305        let producer_ack = producer.clone();
306        if let Some(producer) = producer {
307            self.record_producer_success(
308                stream_id.clone(),
309                producer,
310                ProducerAppendRecord {
311                    start_offset: offset,
312                    next_offset,
313                    closed,
314                },
315                vec![ProducerAppendRecord {
316                    start_offset: offset,
317                    next_offset,
318                    closed,
319                }],
320            );
321        }
322        let object = ObjectPayloadRef {
323            start_offset: offset,
324            end_offset: next_offset,
325            s3_path: payload.s3_path,
326            object_size: payload.object_size,
327        };
328        let slot = self
329            .stream_slot_mut(&stream_id)
330            .expect("stream existence checked before external append mutation");
331        slot.cold.push_external_segment(object.clone());
332        slot.integrity.append_external(
333            &stream_id,
334            object.start_offset,
335            object.end_offset,
336            &object.s3_path,
337            object.object_size,
338        );
339        slot.message_records.push(StreamMessageRecord {
340            start_offset: offset,
341            end_offset: next_offset,
342        });
343        StreamResponse::Appended {
344            offset,
345            next_offset,
346            closed: close_after,
347            deduplicated: false,
348            producer: producer_ack,
349        }
350    }
351
352    pub fn append_batch_borrowed(
353        &mut self,
354        stream_id: BucketStreamId,
355        content_type: Option<&str>,
356        payloads: &[&[u8]],
357        producer: Option<ProducerRequest>,
358        now_ms: u64,
359    ) -> Result<StreamBatchAppend, StreamResponse> {
360        if payloads.is_empty() {
361            return Err(StreamResponse::error(
362                StreamErrorCode::EmptyAppend,
363                "append batch must contain at least one payload",
364            ));
365        }
366        self.validate_stream_scope(&stream_id)?;
367        validate_producer_request(producer.as_ref())?;
368        if self.expire_stream_if_due(&stream_id, now_ms) {
369            return Err(StreamResponse::error(
370                StreamErrorCode::StreamNotFound,
371                format!("stream '{stream_id}' does not exist"),
372            ));
373        }
374        if self
375            .stream_metadata(&stream_id)
376            .is_some_and(is_soft_deleted)
377        {
378            return Err(StreamResponse::error(
379                StreamErrorCode::StreamGone,
380                format!("stream '{stream_id}' is gone"),
381            ));
382        }
383        let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
384        if let ProducerDecision::Duplicate { items, .. } = producer_decision {
385            return Ok(StreamBatchAppend {
386                items: items
387                    .into_iter()
388                    .map(|item| StreamBatchAppendItem {
389                        offset: item.start_offset,
390                        next_offset: item.next_offset,
391                        closed: item.closed,
392                        deduplicated: true,
393                    })
394                    .collect(),
395                deduplicated: true,
396            });
397        }
398
399        let Some(stream) = self.stream_metadata_mut(&stream_id) else {
400            return Err(StreamResponse::error(
401                StreamErrorCode::StreamNotFound,
402                format!("stream '{stream_id}' does not exist"),
403            ));
404        };
405        if stream.status == StreamStatus::Closed {
406            return Err(StreamResponse::error_with_next_offset_and_context(
407                StreamErrorCode::StreamClosed,
408                format!("stream '{stream_id}' is closed"),
409                stream.tail_offset,
410                vec![StreamErrorContext::StreamClosed],
411            ));
412        }
413        let Some(content_type) = content_type else {
414            return Err(StreamResponse::error(
415                StreamErrorCode::MissingContentType,
416                "append batch must include content type",
417            ));
418        };
419        if content_type != stream.content_type {
420            return Err(StreamResponse::error_with_next_offset(
421                StreamErrorCode::ContentTypeMismatch,
422                format!(
423                    "append content type '{content_type}' does not match stream content type '{}'",
424                    stream.content_type
425                ),
426                stream.tail_offset,
427            ));
428        }
429        if payloads.iter().any(|payload| payload.is_empty()) {
430            return Err(StreamResponse::error(
431                StreamErrorCode::EmptyAppend,
432                "append batch payloads must be non-empty",
433            ));
434        }
435
436        let mut items = Vec::with_capacity(payloads.len());
437        for payload in payloads {
438            let offset = stream.tail_offset;
439            let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
440            stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
441            items.push(ProducerAppendRecord {
442                start_offset: offset,
443                next_offset: stream.tail_offset,
444                closed: false,
445            });
446        }
447        let last = items
448            .last()
449            .expect("payloads checked non-empty before append")
450            .clone();
451        renew_stream_ttl(stream, now_ms);
452        self.refresh_ttl_entry(&stream_id);
453        if let Some(producer) = producer {
454            self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
455        }
456        let slot = self
457            .stream_slot_mut(&stream_id)
458            .expect("stream existence checked before batch append mutation");
459        for (item, payload) in items.iter().zip(payloads.iter()) {
460            slot.hot_buffer
461                .push(item.start_offset, item.next_offset, payload);
462        }
463        for (item, payload) in items.iter().zip(payloads.iter()) {
464            slot.integrity
465                .append_payload(&stream_id, item.start_offset, item.next_offset, payload);
466        }
467        slot.message_records
468            .extend(items.iter().map(|item| StreamMessageRecord {
469                start_offset: item.start_offset,
470                end_offset: item.next_offset,
471            }));
472        Ok(StreamBatchAppend {
473            items: items
474                .into_iter()
475                .map(|item| StreamBatchAppendItem {
476                    offset: item.start_offset,
477                    next_offset: item.next_offset,
478                    closed: item.closed,
479                    deduplicated: false,
480                })
481                .collect(),
482            deduplicated: false,
483        })
484    }
485
486    fn evaluate_producer(
487        &self,
488        stream_id: &BucketStreamId,
489        producer: Option<&ProducerRequest>,
490    ) -> Result<ProducerDecision, StreamResponse> {
491        let Some(producer) = producer else {
492            return Ok(ProducerDecision::Accept);
493        };
494        let Some(states) = self.stream_slot(stream_id).map(|slot| &slot.producers) else {
495            return Ok(ProducerDecision::Accept);
496        };
497        let Some(state) = states.get(&producer.producer_id) else {
498            if producer.producer_seq == 0 {
499                return Ok(ProducerDecision::Accept);
500            }
501            return Err(StreamResponse::error_with_context(
502                StreamErrorCode::ProducerSeqConflict,
503                format!(
504                    "producer '{}' expected sequence 0, received {}",
505                    producer.producer_id, producer.producer_seq
506                ),
507                vec![StreamErrorContext::ProducerSeqConflict {
508                    expected_seq: 0,
509                    received_seq: producer.producer_seq,
510                }],
511            ));
512        };
513
514        if producer.producer_epoch < state.producer_epoch {
515            return Err(StreamResponse::error_with_context(
516                StreamErrorCode::ProducerEpochStale,
517                format!(
518                    "producer '{}' epoch {} is stale; current epoch is {}",
519                    producer.producer_id, producer.producer_epoch, state.producer_epoch
520                ),
521                vec![StreamErrorContext::ProducerEpochStale {
522                    current_epoch: state.producer_epoch,
523                }],
524            ));
525        }
526        if producer.producer_epoch > state.producer_epoch {
527            if producer.producer_seq == 0 {
528                return Ok(ProducerDecision::Accept);
529            }
530            return Err(StreamResponse::error(
531                StreamErrorCode::InvalidProducer,
532                format!(
533                    "producer '{}' new epoch {} must start at sequence 0",
534                    producer.producer_id, producer.producer_epoch
535                ),
536            ));
537        }
538
539        if producer.producer_seq <= state.producer_seq {
540            return Ok(ProducerDecision::Duplicate {
541                offset: state.last_start_offset,
542                next_offset: state.last_next_offset,
543                closed: state.last_closed,
544                producer: ProducerRequest {
545                    producer_id: producer.producer_id.clone(),
546                    producer_epoch: state.producer_epoch,
547                    producer_seq: state.producer_seq,
548                },
549                items: state.last_items.clone(),
550            });
551        }
552        if producer.producer_seq == state.producer_seq + 1 {
553            return Ok(ProducerDecision::Accept);
554        }
555        Err(StreamResponse::error_with_context(
556            StreamErrorCode::ProducerSeqConflict,
557            format!(
558                "producer '{}' expected sequence {}, received {}",
559                producer.producer_id,
560                state.producer_seq + 1,
561                producer.producer_seq
562            ),
563            vec![StreamErrorContext::ProducerSeqConflict {
564                expected_seq: state.producer_seq + 1,
565                received_seq: producer.producer_seq,
566            }],
567        ))
568    }
569
570    fn record_producer_success(
571        &mut self,
572        stream_id: BucketStreamId,
573        producer: ProducerRequest,
574        last: ProducerAppendRecord,
575        last_items: Vec<ProducerAppendRecord>,
576    ) {
577        self.stream_slot_mut(&stream_id)
578            .expect("stream existence checked before producer mutation")
579            .producers
580            .insert(producer.producer_id, ProducerState {
581                producer_epoch: producer.producer_epoch,
582                producer_seq: producer.producer_seq,
583                last_start_offset: last.start_offset,
584                last_next_offset: last.next_offset,
585                last_closed: last.closed,
586                last_items,
587            });
588    }
589}
590
591#[derive(Debug, Clone, PartialEq, Eq)]
592enum ProducerDecision {
593    Accept,
594    Duplicate {
595        offset: u64,
596        next_offset: u64,
597        closed: bool,
598        producer: ProducerRequest,
599        items: Vec<ProducerAppendRecord>,
600    },
601}
602
603fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
604    let Some(incoming) = incoming else {
605        return Ok(());
606    };
607    if let Some(last) = stream.last_stream_seq.as_deref()
608        && incoming <= last
609    {
610        return Err(StreamResponse::error_with_next_offset(
611            StreamErrorCode::StreamSeqConflict,
612            format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
613            stream.tail_offset,
614        ));
615    }
616    Ok(())
617}