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