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::StreamMetadata;
15use super::StreamResponse;
16use super::StreamStateMachine;
17use super::StreamStatus;
18use super::canonical_json_record_ends;
19use super::prepare_record_append;
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            record_match,
35        } = input;
36        if let Err(response) = self.validate_stream_scope(&stream_id) {
37            return response;
38        }
39        if let Err(response) = validate_producer_request(producer.as_ref()) {
40            return response;
41        }
42
43        let Some(_) = self.stream_metadata(&stream_id) else {
44            return StreamResponse::error(
45                StreamErrorCode::StreamNotFound,
46                format!("stream '{stream_id}' does not exist"),
47            );
48        };
49        if self.expire_stream_if_due(&stream_id, now_ms) {
50            return StreamResponse::error(
51                StreamErrorCode::StreamNotFound,
52                format!("stream '{stream_id}' does not exist"),
53            );
54        }
55        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
56            Ok(decision) => decision,
57            Err(response) => return response,
58        };
59        if let ProducerDecision::Duplicate {
60            offset,
61            next_offset,
62            closed,
63            producer,
64            ..
65        } = producer_decision
66        {
67            if payload.is_empty() {
68                return StreamResponse::Closed {
69                    next_offset,
70                    deduplicated: true,
71                    producer: Some(producer),
72                };
73            }
74            return StreamResponse::Appended {
75                offset,
76                next_offset,
77                closed,
78                deduplicated: true,
79                producer: Some(producer),
80            };
81        }
82
83        if let Err(response) = self.validate_record_match(&stream_id, record_match) {
84            return response;
85        }
86
87        let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
88        let record_ends = match content_type {
89            Some(value) => match canonical_json_record_ends(value, payload) {
90                Ok(record_ends) => record_ends,
91                Err(_) => {
92                    return StreamResponse::error(
93                        StreamErrorCode::InvalidRecordBoundaries,
94                        "application/json append payload must use canonical newline boundaries",
95                    );
96                }
97            },
98            None => Vec::new(),
99        };
100        let prepared_record_append = {
101            let slot = self
102                .stream_slot(&stream_id)
103                .expect("stream existence checked before record validation");
104            match prepare_record_append(
105                slot.record_index.as_ref(),
106                super::is_json_record_content_type(&slot.metadata.content_type),
107                slot.metadata.tail_offset,
108                payload_len,
109                &record_ends,
110            ) {
111                Ok(prepared) => prepared,
112                Err(response) => return response,
113            }
114        };
115        let record_range = prepared_record_append
116            .as_ref()
117            .map(crate::PreparedRecordAppend::range);
118
119        let Some(stream) = self.stream_metadata_mut(&stream_id) else {
120            unreachable!("stream existence checked before producer evaluation");
121        };
122
123        if stream.status == StreamStatus::Closed {
124            if close_after && payload.is_empty() {
125                return StreamResponse::Closed {
126                    next_offset: stream.tail_offset,
127                    deduplicated: false,
128                    producer: None,
129                };
130            }
131            return StreamResponse::error_with_next_offset_and_context(
132                StreamErrorCode::StreamClosed,
133                format!("stream '{stream_id}' is closed"),
134                stream.tail_offset,
135                vec![StreamErrorContext::StreamClosed],
136            );
137        }
138
139        if payload.is_empty() && !close_after {
140            return StreamResponse::error(
141                StreamErrorCode::EmptyAppend,
142                "append payload must be non-empty unless closing the stream",
143            );
144        }
145
146        if !payload.is_empty() {
147            let Some(content_type) = content_type else {
148                return StreamResponse::error(
149                    StreamErrorCode::MissingContentType,
150                    "append with a body must include content type",
151                );
152            };
153            if content_type != stream.content_type {
154                return StreamResponse::error_with_next_offset(
155                    StreamErrorCode::ContentTypeMismatch,
156                    format!(
157                        "append content type '{content_type}' does not match stream content type '{}'",
158                        stream.content_type
159                    ),
160                    stream.tail_offset,
161                );
162            }
163        }
164
165        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
166            return response;
167        }
168
169        let offset = stream.tail_offset;
170        stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
171        if let Some(seq) = stream_seq {
172            stream.last_stream_seq = Some(seq);
173        }
174        renew_stream_ttl(stream, now_ms);
175        if close_after {
176            stream.status = StreamStatus::Closed;
177        }
178        let closed = stream.status == StreamStatus::Closed;
179        let next_offset = stream.tail_offset;
180        self.refresh_ttl_entry(&stream_id);
181        let producer_ack = producer.clone();
182        if let Some(producer) = producer {
183            self.record_producer_success(
184                stream_id.clone(),
185                producer,
186                ProducerAppendRecord {
187                    start_offset: offset,
188                    next_offset,
189                    closed,
190                    record_start: record_range.map(|range| range.first_record),
191                    record_next: record_range.map(|range| range.next_record),
192                },
193                vec![ProducerAppendRecord {
194                    start_offset: offset,
195                    next_offset,
196                    closed,
197                    record_start: record_range.map(|range| range.first_record),
198                    record_next: record_range.map(|range| range.next_record),
199                }],
200            );
201        }
202
203        if payload.is_empty() {
204            StreamResponse::Closed {
205                next_offset,
206                deduplicated: false,
207                producer: producer_ack,
208            }
209        } else {
210            let slot = self
211                .stream_slot_mut(&stream_id)
212                .expect("stream existence checked before append mutation");
213            if let (Some(index), Some(prepared)) =
214                (slot.record_index.as_mut(), prepared_record_append)
215            {
216                let _range = index.commit_append(prepared);
217            }
218            slot.hot_buffer.push(offset, next_offset, payload);
219            slot.integrity
220                .append_payload(&stream_id, offset, next_offset, payload);
221            slot.message_records
222                .extend(Self::message_records_for_append(
223                    offset,
224                    next_offset,
225                    &record_ends,
226                ));
227            StreamResponse::Appended {
228                offset,
229                next_offset,
230                closed: close_after,
231                deduplicated: false,
232                producer: producer_ack,
233            }
234        }
235    }
236
237    pub(super) fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
238        let AppendExternalInput {
239            stream_id,
240            content_type,
241            payload,
242            record_ends,
243            close_after,
244            stream_seq,
245            producer,
246            now_ms,
247            record_match,
248        } = input;
249        if let Err(response) = validate_external_payload_ref(&payload) {
250            return response;
251        }
252        if let Err(response) = self.validate_stream_scope(&stream_id) {
253            return response;
254        }
255        if let Err(response) = validate_producer_request(producer.as_ref()) {
256            return response;
257        }
258        let Some(_) = self.stream_metadata(&stream_id) else {
259            return StreamResponse::error(
260                StreamErrorCode::StreamNotFound,
261                format!("stream '{stream_id}' does not exist"),
262            );
263        };
264        if self.expire_stream_if_due(&stream_id, now_ms) {
265            return StreamResponse::error(
266                StreamErrorCode::StreamNotFound,
267                format!("stream '{stream_id}' does not exist"),
268            );
269        }
270        let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
271            Ok(decision) => decision,
272            Err(response) => return response,
273        };
274        if let ProducerDecision::Duplicate {
275            offset,
276            next_offset,
277            closed,
278            producer,
279            ..
280        } = producer_decision
281        {
282            return StreamResponse::Appended {
283                offset,
284                next_offset,
285                closed,
286                deduplicated: true,
287                producer: Some(producer),
288            };
289        }
290
291        if let Err(response) = self.validate_record_match(&stream_id, record_match) {
292            return response;
293        }
294
295        let prepared_record_append = {
296            let slot = self
297                .stream_slot(&stream_id)
298                .expect("stream existence checked before record validation");
299            match prepare_record_append(
300                slot.record_index.as_ref(),
301                super::is_json_record_content_type(&slot.metadata.content_type),
302                slot.metadata.tail_offset,
303                payload.payload_len,
304                &record_ends,
305            ) {
306                Ok(prepared) => prepared,
307                Err(response) => return response,
308            }
309        };
310        let record_range = prepared_record_append
311            .as_ref()
312            .map(crate::PreparedRecordAppend::range);
313
314        let Some(stream) = self.stream_metadata(&stream_id) else {
315            unreachable!("stream existence checked before producer evaluation");
316        };
317        if stream.status == StreamStatus::Closed {
318            return StreamResponse::error_with_next_offset_and_context(
319                StreamErrorCode::StreamClosed,
320                format!("stream '{stream_id}' is closed"),
321                stream.tail_offset,
322                vec![StreamErrorContext::StreamClosed],
323            );
324        }
325        let Some(content_type) = content_type else {
326            return StreamResponse::error(
327                StreamErrorCode::MissingContentType,
328                "append with a body must include content type",
329            );
330        };
331        if content_type != stream.content_type {
332            return StreamResponse::error_with_next_offset(
333                StreamErrorCode::ContentTypeMismatch,
334                format!(
335                    "append content type '{content_type}' does not match stream content type '{}'",
336                    stream.content_type
337                ),
338                stream.tail_offset,
339            );
340        }
341        if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
342            return response;
343        }
344        let offset = stream.tail_offset;
345        let next_offset = offset.saturating_add(payload.payload_len);
346        let stream = self
347            .stream_metadata_mut(&stream_id)
348            .expect("stream existence checked before external append mutation");
349        stream.tail_offset = next_offset;
350        if let Some(seq) = stream_seq {
351            stream.last_stream_seq = Some(seq);
352        }
353        renew_stream_ttl(stream, now_ms);
354        if close_after {
355            stream.status = StreamStatus::Closed;
356        }
357        let closed = stream.status == StreamStatus::Closed;
358        self.refresh_ttl_entry(&stream_id);
359        let producer_ack = producer.clone();
360        if let Some(producer) = producer {
361            self.record_producer_success(
362                stream_id.clone(),
363                producer,
364                ProducerAppendRecord {
365                    start_offset: offset,
366                    next_offset,
367                    closed,
368                    record_start: record_range.map(|range| range.first_record),
369                    record_next: record_range.map(|range| range.next_record),
370                },
371                vec![ProducerAppendRecord {
372                    start_offset: offset,
373                    next_offset,
374                    closed,
375                    record_start: record_range.map(|range| range.first_record),
376                    record_next: record_range.map(|range| range.next_record),
377                }],
378            );
379        }
380        let object = ObjectPayloadRef {
381            start_offset: offset,
382            end_offset: next_offset,
383            s3_path: payload.s3_path,
384            object_size: payload.object_size,
385        };
386        let slot = self
387            .stream_slot_mut(&stream_id)
388            .expect("stream existence checked before external append mutation");
389        if let (Some(index), Some(prepared)) = (slot.record_index.as_mut(), prepared_record_append)
390        {
391            let _range = index.commit_append(prepared);
392        }
393        slot.cold.push_external_segment(object.clone());
394        slot.integrity.append_external(
395            &stream_id,
396            object.start_offset,
397            object.end_offset,
398            &object.s3_path,
399            object.object_size,
400        );
401        slot.message_records
402            .extend(Self::message_records_for_append(
403                offset,
404                next_offset,
405                &record_ends,
406            ));
407        StreamResponse::Appended {
408            offset,
409            next_offset,
410            closed: close_after,
411            deduplicated: false,
412            producer: producer_ack,
413        }
414    }
415
416    pub fn append_batch_borrowed(
417        &mut self,
418        stream_id: BucketStreamId,
419        content_type: Option<&str>,
420        payloads: &[&[u8]],
421        producer: Option<ProducerRequest>,
422        now_ms: u64,
423    ) -> Result<StreamBatchAppend, StreamResponse> {
424        if payloads.is_empty() {
425            return Err(StreamResponse::error(
426                StreamErrorCode::EmptyAppend,
427                "append batch must contain at least one payload",
428            ));
429        }
430        self.validate_stream_scope(&stream_id)?;
431        validate_producer_request(producer.as_ref())?;
432        if self.expire_stream_if_due(&stream_id, now_ms) {
433            return Err(StreamResponse::error(
434                StreamErrorCode::StreamNotFound,
435                format!("stream '{stream_id}' does not exist"),
436            ));
437        }
438        let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
439        if let ProducerDecision::Duplicate { items, .. } = producer_decision {
440            return Ok(StreamBatchAppend {
441                items: items
442                    .into_iter()
443                    .map(|item| StreamBatchAppendItem {
444                        offset: item.start_offset,
445                        next_offset: item.next_offset,
446                        closed: item.closed,
447                        deduplicated: true,
448                    })
449                    .collect(),
450                deduplicated: true,
451            });
452        }
453
454        let Some(stream) = self.stream_metadata(&stream_id) else {
455            return Err(StreamResponse::error(
456                StreamErrorCode::StreamNotFound,
457                format!("stream '{stream_id}' does not exist"),
458            ));
459        };
460        if stream.status == StreamStatus::Closed {
461            return Err(StreamResponse::error_with_next_offset_and_context(
462                StreamErrorCode::StreamClosed,
463                format!("stream '{stream_id}' is closed"),
464                stream.tail_offset,
465                vec![StreamErrorContext::StreamClosed],
466            ));
467        }
468        let Some(content_type) = content_type else {
469            return Err(StreamResponse::error(
470                StreamErrorCode::MissingContentType,
471                "append batch must include content type",
472            ));
473        };
474        if content_type != stream.content_type {
475            return Err(StreamResponse::error_with_next_offset(
476                StreamErrorCode::ContentTypeMismatch,
477                format!(
478                    "append content type '{content_type}' does not match stream content type '{}'",
479                    stream.content_type
480                ),
481                stream.tail_offset,
482            ));
483        }
484        if payloads.iter().any(|payload| payload.is_empty()) {
485            return Err(StreamResponse::error(
486                StreamErrorCode::EmptyAppend,
487                "append batch payloads must be non-empty",
488            ));
489        }
490
491        let base_offset = stream.tail_offset;
492        let mut total_payload_len = 0_u64;
493        let mut combined_record_ends = Vec::new();
494        let mut record_counts = Vec::with_capacity(payloads.len());
495        let mut all_record_ends = Vec::with_capacity(payloads.len());
496        for payload in payloads {
497            let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
498            let record_ends = canonical_json_record_ends(content_type, payload).map_err(|_| {
499                StreamResponse::error(
500                    StreamErrorCode::InvalidRecordBoundaries,
501                    "application/json append payload must use canonical newline boundaries",
502                )
503            })?;
504            record_counts.push(u64::try_from(record_ends.len()).map_err(|_| {
505                StreamResponse::error(
506                    StreamErrorCode::InvalidRecordBoundaries,
507                    "record count exceeds the supported range",
508                )
509            })?);
510            for end in &record_ends {
511                combined_record_ends.push(total_payload_len.checked_add(*end).ok_or_else(
512                    || {
513                        StreamResponse::error(
514                            StreamErrorCode::InvalidRecordBoundaries,
515                            "append batch payload length exceeds the supported range",
516                        )
517                    },
518                )?);
519            }
520            total_payload_len = total_payload_len.checked_add(payload_len).ok_or_else(|| {
521                StreamResponse::error(
522                    StreamErrorCode::InvalidRecordBoundaries,
523                    "append batch payload length exceeds the supported range",
524                )
525            })?;
526            all_record_ends.push(record_ends);
527        }
528        let prepared_record_append = {
529            let slot = self
530                .stream_slot(&stream_id)
531                .expect("stream existence checked before record validation");
532            prepare_record_append(
533                slot.record_index.as_ref(),
534                super::is_json_record_content_type(content_type),
535                base_offset,
536                total_payload_len,
537                &combined_record_ends,
538            )?
539        };
540        let mut next_record = prepared_record_append
541            .as_ref()
542            .map(crate::PreparedRecordAppend::range)
543            .map(|range| range.first_record);
544        let mut item_record_ranges = Vec::with_capacity(record_counts.len());
545        for count in record_counts {
546            let range = match next_record {
547                Some(first_record) => {
548                    let next = first_record.checked_add(count).ok_or_else(|| {
549                        StreamResponse::error(
550                            StreamErrorCode::InvalidRecordBoundaries,
551                            "append batch record count exceeds the supported range",
552                        )
553                    })?;
554                    next_record = Some(next);
555                    Some(crate::StreamRecordRange {
556                        first_record,
557                        next_record: next,
558                    })
559                }
560                None => None,
561            };
562            item_record_ranges.push(range);
563        }
564
565        let stream = self
566            .stream_metadata_mut(&stream_id)
567            .expect("stream existence checked before batch append mutation");
568
569        let mut items = Vec::with_capacity(payloads.len());
570        for (payload, record_range) in payloads.iter().zip(item_record_ranges) {
571            let offset = stream.tail_offset;
572            let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
573            stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
574            items.push(ProducerAppendRecord {
575                start_offset: offset,
576                next_offset: stream.tail_offset,
577                closed: false,
578                record_start: record_range.map(|range| range.first_record),
579                record_next: record_range.map(|range| range.next_record),
580            });
581        }
582        let last = items
583            .last()
584            .expect("payloads checked non-empty before append")
585            .clone();
586        renew_stream_ttl(stream, now_ms);
587        self.refresh_ttl_entry(&stream_id);
588        if let Some(producer) = producer {
589            self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
590        }
591        let slot = self
592            .stream_slot_mut(&stream_id)
593            .expect("stream existence checked before batch append mutation");
594        if let (Some(index), Some(prepared)) = (slot.record_index.as_mut(), prepared_record_append)
595        {
596            let _range = index.commit_append(prepared);
597        }
598        for (item, payload) in items.iter().zip(payloads.iter()) {
599            slot.hot_buffer
600                .push(item.start_offset, item.next_offset, payload);
601        }
602        for (item, payload) in items.iter().zip(payloads.iter()) {
603            slot.integrity
604                .append_payload(&stream_id, item.start_offset, item.next_offset, payload);
605        }
606        for (item, record_ends) in items.iter().zip(all_record_ends.iter()) {
607            slot.message_records
608                .extend(Self::message_records_for_append(
609                    item.start_offset,
610                    item.next_offset,
611                    record_ends,
612                ));
613        }
614        Ok(StreamBatchAppend {
615            items: items
616                .into_iter()
617                .map(|item| StreamBatchAppendItem {
618                    offset: item.start_offset,
619                    next_offset: item.next_offset,
620                    closed: item.closed,
621                    deduplicated: false,
622                })
623                .collect(),
624            deduplicated: false,
625        })
626    }
627
628    fn validate_record_match(
629        &self,
630        stream_id: &BucketStreamId,
631        expected: Option<u64>,
632    ) -> Result<(), StreamResponse> {
633        let Some(expected) = expected else {
634            return Ok(());
635        };
636        let Some(slot) = self.stream_slot(stream_id) else {
637            return Ok(());
638        };
639        let Some(index) = slot.record_index.as_ref() else {
640            return Err(StreamResponse::error(
641                StreamErrorCode::InvalidRecordBoundaries,
642                "Stream-Record-Match requires active JSON record coordinates",
643            ));
644        };
645        let current = index
646            .range()
647            .map_err(|_| {
648                StreamResponse::error(
649                    StreamErrorCode::InvalidRecordBoundaries,
650                    "stream record index is invalid",
651                )
652            })?
653            .next_record;
654        if current == expected {
655            return Ok(());
656        }
657        Err(StreamResponse::error_with_next_offset_and_context(
658            StreamErrorCode::RecordPreconditionFailed,
659            format!("record tail is {current}, expected {expected}"),
660            slot.metadata.tail_offset,
661            vec![StreamErrorContext::RecordTailMismatch {
662                current_record: current,
663            }],
664        ))
665    }
666
667    fn evaluate_producer(
668        &self,
669        stream_id: &BucketStreamId,
670        producer: Option<&ProducerRequest>,
671    ) -> Result<ProducerDecision, StreamResponse> {
672        let Some(producer) = producer else {
673            return Ok(ProducerDecision::Accept);
674        };
675        let Some(states) = self.stream_slot(stream_id).map(|slot| &slot.producers) else {
676            return Ok(ProducerDecision::Accept);
677        };
678        let Some(state) = states.get(&producer.producer_id) else {
679            if producer.producer_seq == 0 {
680                return Ok(ProducerDecision::Accept);
681            }
682            return Err(StreamResponse::error_with_context(
683                StreamErrorCode::ProducerSeqConflict,
684                format!(
685                    "producer '{}' expected sequence 0, received {}",
686                    producer.producer_id, producer.producer_seq
687                ),
688                vec![StreamErrorContext::ProducerSeqConflict {
689                    expected_seq: 0,
690                    received_seq: producer.producer_seq,
691                }],
692            ));
693        };
694
695        if producer.producer_epoch < state.producer_epoch {
696            return Err(StreamResponse::error_with_context(
697                StreamErrorCode::ProducerEpochStale,
698                format!(
699                    "producer '{}' epoch {} is stale; current epoch is {}",
700                    producer.producer_id, producer.producer_epoch, state.producer_epoch
701                ),
702                vec![StreamErrorContext::ProducerEpochStale {
703                    current_epoch: state.producer_epoch,
704                }],
705            ));
706        }
707        if producer.producer_epoch > state.producer_epoch {
708            if producer.producer_seq == 0 {
709                return Ok(ProducerDecision::Accept);
710            }
711            return Err(StreamResponse::error(
712                StreamErrorCode::InvalidProducer,
713                format!(
714                    "producer '{}' new epoch {} must start at sequence 0",
715                    producer.producer_id, producer.producer_epoch
716                ),
717            ));
718        }
719
720        if producer.producer_seq <= state.producer_seq {
721            return Ok(ProducerDecision::Duplicate {
722                offset: state.last_start_offset,
723                next_offset: state.last_next_offset,
724                closed: state.last_closed,
725                producer: ProducerRequest {
726                    producer_id: producer.producer_id.clone(),
727                    producer_epoch: state.producer_epoch,
728                    producer_seq: state.producer_seq,
729                },
730                items: state.last_items.clone(),
731            });
732        }
733        if producer.producer_seq == state.producer_seq + 1 {
734            return Ok(ProducerDecision::Accept);
735        }
736        Err(StreamResponse::error_with_context(
737            StreamErrorCode::ProducerSeqConflict,
738            format!(
739                "producer '{}' expected sequence {}, received {}",
740                producer.producer_id,
741                state.producer_seq + 1,
742                producer.producer_seq
743            ),
744            vec![StreamErrorContext::ProducerSeqConflict {
745                expected_seq: state.producer_seq + 1,
746                received_seq: producer.producer_seq,
747            }],
748        ))
749    }
750
751    fn record_producer_success(
752        &mut self,
753        stream_id: BucketStreamId,
754        producer: ProducerRequest,
755        last: ProducerAppendRecord,
756        last_items: Vec<ProducerAppendRecord>,
757    ) {
758        self.stream_slot_mut(&stream_id)
759            .expect("stream existence checked before producer mutation")
760            .producers
761            .insert(producer.producer_id, ProducerState {
762                producer_epoch: producer.producer_epoch,
763                producer_seq: producer.producer_seq,
764                last_start_offset: last.start_offset,
765                last_next_offset: last.next_offset,
766                last_closed: last.closed,
767                last_items,
768            });
769    }
770}
771
772#[derive(Debug, Clone, PartialEq, Eq)]
773enum ProducerDecision {
774    Accept,
775    Duplicate {
776        offset: u64,
777        next_offset: u64,
778        closed: bool,
779        producer: ProducerRequest,
780        items: Vec<ProducerAppendRecord>,
781    },
782}
783
784fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
785    let Some(incoming) = incoming else {
786        return Ok(());
787    };
788    if let Some(last) = stream.last_stream_seq.as_deref()
789        && incoming <= last
790    {
791        return Err(StreamResponse::error_with_next_offset(
792            StreamErrorCode::StreamSeqConflict,
793            format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
794            stream.tail_offset,
795        ));
796    }
797    Ok(())
798}