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