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