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::StreamMessageRecord;
15use super::StreamMetadata;
16use super::StreamResponse;
17use super::StreamStateMachine;
18use super::StreamStatus;
19use super::is_soft_deleted;
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 } = input;
35 if let Err(response) = self.validate_stream_scope(&stream_id) {
36 return response;
37 }
38 if let Err(response) = validate_producer_request(producer.as_ref()) {
39 return response;
40 }
41
42 let Some(_) = self.stream_metadata(&stream_id) else {
43 return StreamResponse::error(
44 StreamErrorCode::StreamNotFound,
45 format!("stream '{stream_id}' does not exist"),
46 );
47 };
48 if self.expire_stream_if_due(&stream_id, now_ms) {
49 return StreamResponse::error(
50 StreamErrorCode::StreamNotFound,
51 format!("stream '{stream_id}' does not exist"),
52 );
53 }
54 if self
55 .stream_metadata(&stream_id)
56 .is_some_and(is_soft_deleted)
57 {
58 return StreamResponse::error(
59 StreamErrorCode::StreamGone,
60 format!("stream '{stream_id}' is gone"),
61 );
62 }
63 let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
64 Ok(decision) => decision,
65 Err(response) => return response,
66 };
67 if let ProducerDecision::Duplicate {
68 offset,
69 next_offset,
70 closed,
71 producer,
72 ..
73 } = producer_decision
74 {
75 if payload.is_empty() {
76 return StreamResponse::Closed {
77 next_offset,
78 deduplicated: true,
79 producer: Some(producer),
80 };
81 }
82 return StreamResponse::Appended {
83 offset,
84 next_offset,
85 closed,
86 deduplicated: true,
87 producer: Some(producer),
88 };
89 }
90
91 let Some(stream) = self.stream_metadata_mut(&stream_id) else {
92 unreachable!("stream existence checked before producer evaluation");
93 };
94
95 if stream.status == StreamStatus::Closed {
96 if close_after && payload.is_empty() {
97 return StreamResponse::Closed {
98 next_offset: stream.tail_offset,
99 deduplicated: false,
100 producer: None,
101 };
102 }
103 return StreamResponse::error_with_next_offset_and_context(
104 StreamErrorCode::StreamClosed,
105 format!("stream '{stream_id}' is closed"),
106 stream.tail_offset,
107 vec![StreamErrorContext::StreamClosed],
108 );
109 }
110
111 if payload.is_empty() && !close_after {
112 return StreamResponse::error(
113 StreamErrorCode::EmptyAppend,
114 "append payload must be non-empty unless closing the stream",
115 );
116 }
117
118 if !payload.is_empty() {
119 let Some(content_type) = content_type else {
120 return StreamResponse::error(
121 StreamErrorCode::MissingContentType,
122 "append with a body must include content type",
123 );
124 };
125 if content_type != stream.content_type {
126 return StreamResponse::error_with_next_offset(
127 StreamErrorCode::ContentTypeMismatch,
128 format!(
129 "append content type '{content_type}' does not match stream content type '{}'",
130 stream.content_type
131 ),
132 stream.tail_offset,
133 );
134 }
135 }
136
137 if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
138 return response;
139 }
140
141 let offset = stream.tail_offset;
142 let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
143 stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
144 if let Some(seq) = stream_seq {
145 stream.last_stream_seq = Some(seq);
146 }
147 renew_stream_ttl(stream, now_ms);
148 if close_after {
149 stream.status = StreamStatus::Closed;
150 }
151 let closed = stream.status == StreamStatus::Closed;
152 let next_offset = stream.tail_offset;
153 self.refresh_ttl_entry(&stream_id);
154 let producer_ack = producer.clone();
155 if let Some(producer) = producer {
156 self.record_producer_success(
157 stream_id.clone(),
158 producer,
159 ProducerAppendRecord {
160 start_offset: offset,
161 next_offset,
162 closed,
163 },
164 vec![ProducerAppendRecord {
165 start_offset: offset,
166 next_offset,
167 closed,
168 }],
169 );
170 }
171
172 if payload.is_empty() {
173 StreamResponse::Closed {
174 next_offset,
175 deduplicated: false,
176 producer: producer_ack,
177 }
178 } else {
179 let slot = self
180 .stream_slot_mut(&stream_id)
181 .expect("stream existence checked before append mutation");
182 slot.hot_buffer.push(offset, next_offset, payload);
183 slot.integrity
184 .append_payload(&stream_id, offset, next_offset, payload);
185 slot.message_records.push(StreamMessageRecord {
186 start_offset: offset,
187 end_offset: next_offset,
188 });
189 StreamResponse::Appended {
190 offset,
191 next_offset,
192 closed: close_after,
193 deduplicated: false,
194 producer: producer_ack,
195 }
196 }
197 }
198
199 pub(super) fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
200 let AppendExternalInput {
201 stream_id,
202 content_type,
203 payload,
204 close_after,
205 stream_seq,
206 producer,
207 now_ms,
208 } = input;
209 if let Err(response) = validate_external_payload_ref(&payload) {
210 return response;
211 }
212 if let Err(response) = self.validate_stream_scope(&stream_id) {
213 return response;
214 }
215 if let Err(response) = validate_producer_request(producer.as_ref()) {
216 return response;
217 }
218 let Some(_) = self.stream_metadata(&stream_id) else {
219 return StreamResponse::error(
220 StreamErrorCode::StreamNotFound,
221 format!("stream '{stream_id}' does not exist"),
222 );
223 };
224 if self.expire_stream_if_due(&stream_id, now_ms) {
225 return StreamResponse::error(
226 StreamErrorCode::StreamNotFound,
227 format!("stream '{stream_id}' does not exist"),
228 );
229 }
230 if self
231 .stream_metadata(&stream_id)
232 .is_some_and(is_soft_deleted)
233 {
234 return StreamResponse::error(
235 StreamErrorCode::StreamGone,
236 format!("stream '{stream_id}' is gone"),
237 );
238 }
239 let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
240 Ok(decision) => decision,
241 Err(response) => return response,
242 };
243 if let ProducerDecision::Duplicate {
244 offset,
245 next_offset,
246 closed,
247 producer,
248 ..
249 } = producer_decision
250 {
251 return StreamResponse::Appended {
252 offset,
253 next_offset,
254 closed,
255 deduplicated: true,
256 producer: Some(producer),
257 };
258 }
259
260 let Some(stream) = self.stream_metadata(&stream_id) else {
261 unreachable!("stream existence checked before producer evaluation");
262 };
263 if stream.status == StreamStatus::Closed {
264 return StreamResponse::error_with_next_offset_and_context(
265 StreamErrorCode::StreamClosed,
266 format!("stream '{stream_id}' is closed"),
267 stream.tail_offset,
268 vec![StreamErrorContext::StreamClosed],
269 );
270 }
271 let Some(content_type) = content_type else {
272 return StreamResponse::error(
273 StreamErrorCode::MissingContentType,
274 "append with a body must include content type",
275 );
276 };
277 if content_type != stream.content_type {
278 return StreamResponse::error_with_next_offset(
279 StreamErrorCode::ContentTypeMismatch,
280 format!(
281 "append content type '{content_type}' does not match stream content type '{}'",
282 stream.content_type
283 ),
284 stream.tail_offset,
285 );
286 }
287 if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
288 return response;
289 }
290 let offset = stream.tail_offset;
291 let next_offset = offset.saturating_add(payload.payload_len);
292 let stream = self
293 .stream_metadata_mut(&stream_id)
294 .expect("stream existence checked before external append mutation");
295 stream.tail_offset = next_offset;
296 if let Some(seq) = stream_seq {
297 stream.last_stream_seq = Some(seq);
298 }
299 renew_stream_ttl(stream, now_ms);
300 if close_after {
301 stream.status = StreamStatus::Closed;
302 }
303 let closed = stream.status == StreamStatus::Closed;
304 self.refresh_ttl_entry(&stream_id);
305 let producer_ack = producer.clone();
306 if let Some(producer) = producer {
307 self.record_producer_success(
308 stream_id.clone(),
309 producer,
310 ProducerAppendRecord {
311 start_offset: offset,
312 next_offset,
313 closed,
314 },
315 vec![ProducerAppendRecord {
316 start_offset: offset,
317 next_offset,
318 closed,
319 }],
320 );
321 }
322 let object = ObjectPayloadRef {
323 start_offset: offset,
324 end_offset: next_offset,
325 s3_path: payload.s3_path,
326 object_size: payload.object_size,
327 };
328 let slot = self
329 .stream_slot_mut(&stream_id)
330 .expect("stream existence checked before external append mutation");
331 slot.cold.push_external_segment(object.clone());
332 slot.integrity.append_external(
333 &stream_id,
334 object.start_offset,
335 object.end_offset,
336 &object.s3_path,
337 object.object_size,
338 );
339 slot.message_records.push(StreamMessageRecord {
340 start_offset: offset,
341 end_offset: next_offset,
342 });
343 StreamResponse::Appended {
344 offset,
345 next_offset,
346 closed: close_after,
347 deduplicated: false,
348 producer: producer_ack,
349 }
350 }
351
352 pub fn append_batch_borrowed(
353 &mut self,
354 stream_id: BucketStreamId,
355 content_type: Option<&str>,
356 payloads: &[&[u8]],
357 producer: Option<ProducerRequest>,
358 now_ms: u64,
359 ) -> Result<StreamBatchAppend, StreamResponse> {
360 if payloads.is_empty() {
361 return Err(StreamResponse::error(
362 StreamErrorCode::EmptyAppend,
363 "append batch must contain at least one payload",
364 ));
365 }
366 self.validate_stream_scope(&stream_id)?;
367 validate_producer_request(producer.as_ref())?;
368 if self.expire_stream_if_due(&stream_id, now_ms) {
369 return Err(StreamResponse::error(
370 StreamErrorCode::StreamNotFound,
371 format!("stream '{stream_id}' does not exist"),
372 ));
373 }
374 if self
375 .stream_metadata(&stream_id)
376 .is_some_and(is_soft_deleted)
377 {
378 return Err(StreamResponse::error(
379 StreamErrorCode::StreamGone,
380 format!("stream '{stream_id}' is gone"),
381 ));
382 }
383 let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
384 if let ProducerDecision::Duplicate { items, .. } = producer_decision {
385 return Ok(StreamBatchAppend {
386 items: items
387 .into_iter()
388 .map(|item| StreamBatchAppendItem {
389 offset: item.start_offset,
390 next_offset: item.next_offset,
391 closed: item.closed,
392 deduplicated: true,
393 })
394 .collect(),
395 deduplicated: true,
396 });
397 }
398
399 let Some(stream) = self.stream_metadata_mut(&stream_id) else {
400 return Err(StreamResponse::error(
401 StreamErrorCode::StreamNotFound,
402 format!("stream '{stream_id}' does not exist"),
403 ));
404 };
405 if stream.status == StreamStatus::Closed {
406 return Err(StreamResponse::error_with_next_offset_and_context(
407 StreamErrorCode::StreamClosed,
408 format!("stream '{stream_id}' is closed"),
409 stream.tail_offset,
410 vec![StreamErrorContext::StreamClosed],
411 ));
412 }
413 let Some(content_type) = content_type else {
414 return Err(StreamResponse::error(
415 StreamErrorCode::MissingContentType,
416 "append batch must include content type",
417 ));
418 };
419 if content_type != stream.content_type {
420 return Err(StreamResponse::error_with_next_offset(
421 StreamErrorCode::ContentTypeMismatch,
422 format!(
423 "append content type '{content_type}' does not match stream content type '{}'",
424 stream.content_type
425 ),
426 stream.tail_offset,
427 ));
428 }
429 if payloads.iter().any(|payload| payload.is_empty()) {
430 return Err(StreamResponse::error(
431 StreamErrorCode::EmptyAppend,
432 "append batch payloads must be non-empty",
433 ));
434 }
435
436 let mut items = Vec::with_capacity(payloads.len());
437 for payload in payloads {
438 let offset = stream.tail_offset;
439 let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
440 stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
441 items.push(ProducerAppendRecord {
442 start_offset: offset,
443 next_offset: stream.tail_offset,
444 closed: false,
445 });
446 }
447 let last = items
448 .last()
449 .expect("payloads checked non-empty before append")
450 .clone();
451 renew_stream_ttl(stream, now_ms);
452 self.refresh_ttl_entry(&stream_id);
453 if let Some(producer) = producer {
454 self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
455 }
456 let slot = self
457 .stream_slot_mut(&stream_id)
458 .expect("stream existence checked before batch append mutation");
459 for (item, payload) in items.iter().zip(payloads.iter()) {
460 slot.hot_buffer
461 .push(item.start_offset, item.next_offset, payload);
462 }
463 for (item, payload) in items.iter().zip(payloads.iter()) {
464 slot.integrity
465 .append_payload(&stream_id, item.start_offset, item.next_offset, payload);
466 }
467 slot.message_records
468 .extend(items.iter().map(|item| StreamMessageRecord {
469 start_offset: item.start_offset,
470 end_offset: item.next_offset,
471 }));
472 Ok(StreamBatchAppend {
473 items: items
474 .into_iter()
475 .map(|item| StreamBatchAppendItem {
476 offset: item.start_offset,
477 next_offset: item.next_offset,
478 closed: item.closed,
479 deduplicated: false,
480 })
481 .collect(),
482 deduplicated: false,
483 })
484 }
485
486 fn evaluate_producer(
487 &self,
488 stream_id: &BucketStreamId,
489 producer: Option<&ProducerRequest>,
490 ) -> Result<ProducerDecision, StreamResponse> {
491 let Some(producer) = producer else {
492 return Ok(ProducerDecision::Accept);
493 };
494 let Some(states) = self.stream_slot(stream_id).map(|slot| &slot.producers) else {
495 return Ok(ProducerDecision::Accept);
496 };
497 let Some(state) = states.get(&producer.producer_id) else {
498 if producer.producer_seq == 0 {
499 return Ok(ProducerDecision::Accept);
500 }
501 return Err(StreamResponse::error_with_context(
502 StreamErrorCode::ProducerSeqConflict,
503 format!(
504 "producer '{}' expected sequence 0, received {}",
505 producer.producer_id, producer.producer_seq
506 ),
507 vec![StreamErrorContext::ProducerSeqConflict {
508 expected_seq: 0,
509 received_seq: producer.producer_seq,
510 }],
511 ));
512 };
513
514 if producer.producer_epoch < state.producer_epoch {
515 return Err(StreamResponse::error_with_context(
516 StreamErrorCode::ProducerEpochStale,
517 format!(
518 "producer '{}' epoch {} is stale; current epoch is {}",
519 producer.producer_id, producer.producer_epoch, state.producer_epoch
520 ),
521 vec![StreamErrorContext::ProducerEpochStale {
522 current_epoch: state.producer_epoch,
523 }],
524 ));
525 }
526 if producer.producer_epoch > state.producer_epoch {
527 if producer.producer_seq == 0 {
528 return Ok(ProducerDecision::Accept);
529 }
530 return Err(StreamResponse::error(
531 StreamErrorCode::InvalidProducer,
532 format!(
533 "producer '{}' new epoch {} must start at sequence 0",
534 producer.producer_id, producer.producer_epoch
535 ),
536 ));
537 }
538
539 if producer.producer_seq <= state.producer_seq {
540 return Ok(ProducerDecision::Duplicate {
541 offset: state.last_start_offset,
542 next_offset: state.last_next_offset,
543 closed: state.last_closed,
544 producer: ProducerRequest {
545 producer_id: producer.producer_id.clone(),
546 producer_epoch: state.producer_epoch,
547 producer_seq: state.producer_seq,
548 },
549 items: state.last_items.clone(),
550 });
551 }
552 if producer.producer_seq == state.producer_seq + 1 {
553 return Ok(ProducerDecision::Accept);
554 }
555 Err(StreamResponse::error_with_context(
556 StreamErrorCode::ProducerSeqConflict,
557 format!(
558 "producer '{}' expected sequence {}, received {}",
559 producer.producer_id,
560 state.producer_seq + 1,
561 producer.producer_seq
562 ),
563 vec![StreamErrorContext::ProducerSeqConflict {
564 expected_seq: state.producer_seq + 1,
565 received_seq: producer.producer_seq,
566 }],
567 ))
568 }
569
570 fn record_producer_success(
571 &mut self,
572 stream_id: BucketStreamId,
573 producer: ProducerRequest,
574 last: ProducerAppendRecord,
575 last_items: Vec<ProducerAppendRecord>,
576 ) {
577 self.stream_slot_mut(&stream_id)
578 .expect("stream existence checked before producer mutation")
579 .producers
580 .insert(producer.producer_id, ProducerState {
581 producer_epoch: producer.producer_epoch,
582 producer_seq: producer.producer_seq,
583 last_start_offset: last.start_offset,
584 last_next_offset: last.next_offset,
585 last_closed: last.closed,
586 last_items,
587 });
588 }
589}
590
591#[derive(Debug, Clone, PartialEq, Eq)]
592enum ProducerDecision {
593 Accept,
594 Duplicate {
595 offset: u64,
596 next_offset: u64,
597 closed: bool,
598 producer: ProducerRequest,
599 items: Vec<ProducerAppendRecord>,
600 },
601}
602
603fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
604 let Some(incoming) = incoming else {
605 return Ok(());
606 };
607 if let Some(last) = stream.last_stream_seq.as_deref()
608 && incoming <= last
609 {
610 return Err(StreamResponse::error_with_next_offset(
611 StreamErrorCode::StreamSeqConflict,
612 format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
613 stream.tail_offset,
614 ));
615 }
616 Ok(())
617}