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