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