use super::AppendExternalInput;
use super::AppendStreamInput;
use super::BucketStreamId;
use super::BucketUsage;
use super::HashMap;
use super::ObjectPayloadRef;
use super::ProducerAppendRecord;
use super::ProducerReceipt;
use super::ProducerRequest;
use super::ProducerState;
use super::StreamBatchAppend;
use super::StreamBatchAppendItem;
use super::StreamCommand;
use super::StreamErrorCode;
use super::StreamErrorContext;
use super::StreamIntegrity;
use super::StreamMetadata;
use super::StreamResponse;
use super::StreamStateMachine;
use super::StreamStatus;
use super::canonical_json_record_ends;
use super::prepare_record_append;
use super::renew_stream_ttl;
use super::validate_external_payload_ref;
use super::validate_producer_request;
struct StreamAppendUndo {
metadata: StreamMetadata,
hot_checkpoint: usize,
message_records_len: usize,
record_checkpoint: Option<usize>,
integrity: StreamIntegrity,
producers: HashMap<String, Option<ProducerState>>,
}
impl StreamStateMachine {
pub fn append_transaction(
&mut self,
commands: Vec<StreamCommand>,
) -> Result<Vec<StreamResponse>, StreamResponse> {
let Some(StreamCommand::Append {
stream_id: first_stream,
..
}) = commands.first()
else {
return Err(StreamResponse::error(
StreamErrorCode::InvalidStreamId,
"append transaction must contain at least one append command",
));
};
let Some(first_affinity) = first_stream.affinity_key.as_deref() else {
return Err(StreamResponse::error(
StreamErrorCode::InvalidStreamId,
"append transaction streams must use path affinity",
));
};
let mut stream_undo = HashMap::<BucketStreamId, StreamAppendUndo>::new();
let mut bucket_undo = HashMap::<String, Option<BucketUsage>>::new();
for command in &commands {
let StreamCommand::Append {
stream_id,
producer,
now_ms,
..
} = command
else {
return Err(StreamResponse::error(
StreamErrorCode::InvalidStreamId,
"append transaction contains a non-append command",
));
};
if stream_id.bucket_id != first_stream.bucket_id
|| stream_id.affinity_key.as_deref() != Some(first_affinity)
{
return Err(StreamResponse::error(
StreamErrorCode::InvalidStreamId,
"append transaction streams must share one bucket and affinity key",
));
}
let Some(slot) = self.stream_slot(stream_id) else {
return Err(StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
));
};
if super::stream_is_expired(&slot.metadata, *now_ms) {
return Err(StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
));
}
bucket_undo
.entry(stream_id.bucket_id.clone())
.or_insert_with(|| self.bucket_usage.get(&stream_id.bucket_id).copied());
let undo = stream_undo
.entry(stream_id.clone())
.or_insert_with(|| StreamAppendUndo {
metadata: slot.metadata.clone(),
hot_checkpoint: slot.hot_buffer.append_checkpoint(),
message_records_len: slot.message_records.len(),
record_checkpoint: slot
.record_index
.as_ref()
.map(crate::StreamRecordIndex::append_checkpoint),
integrity: slot.integrity.clone(),
producers: HashMap::new(),
});
if let Some(producer) = producer {
undo.producers
.entry(producer.producer_id.clone())
.or_insert_with(|| slot.producers.get(&producer.producer_id).cloned());
}
}
let hot_payload_bytes = self.hot_payload_bytes;
let mut responses = Vec::with_capacity(commands.len());
for command in commands {
let StreamCommand::Append {
stream_id,
content_type,
payload,
close_after,
stream_seq,
producer,
now_ms,
record_match,
} = command
else {
unreachable!("transaction commands validated before mutation");
};
let response = self.append_borrowed(AppendStreamInput {
stream_id,
content_type: content_type.as_deref(),
payload: &payload,
close_after,
stream_seq,
producer,
now_ms,
record_match,
});
if matches!(response, StreamResponse::Error { .. }) {
self.rollback_append_transaction(stream_undo, bucket_undo, hot_payload_bytes);
return Err(response);
}
responses.push(response);
}
Ok(responses)
}
fn rollback_append_transaction(
&mut self,
stream_undo: HashMap<BucketStreamId, StreamAppendUndo>,
bucket_undo: HashMap<String, Option<BucketUsage>>,
hot_payload_bytes: u64,
) {
self.hot_payload_bytes = hot_payload_bytes;
for (bucket_id, usage) in bucket_undo {
match usage {
Some(usage) => {
self.bucket_usage.insert(bucket_id, usage);
}
None => {
self.bucket_usage.remove(&bucket_id);
}
}
}
for (stream_id, undo) in stream_undo {
let Some(slot) = self.stream_slot_mut(&stream_id) else {
continue;
};
slot.metadata = undo.metadata;
slot.hot_buffer.rollback_appends(undo.hot_checkpoint);
slot.message_records.truncate(undo.message_records_len);
if let (Some(index), Some(checkpoint)) =
(slot.record_index.as_mut(), undo.record_checkpoint)
{
index.rollback_appends(checkpoint);
}
slot.integrity = undo.integrity;
for (producer_id, producer) in undo.producers {
match producer {
Some(producer) => {
slot.producers.insert(producer_id, producer);
}
None => {
slot.producers.remove(&producer_id);
}
}
}
}
}
pub fn append_borrowed(&mut self, input: AppendStreamInput<'_>) -> StreamResponse {
let AppendStreamInput {
stream_id,
content_type,
payload,
close_after,
stream_seq,
producer,
now_ms,
record_match,
} = input;
if let Err(response) = self.validate_stream_scope(&stream_id) {
return response;
}
if let Err(response) = validate_producer_request(producer.as_ref()) {
return response;
}
let Some(_) = self.stream_metadata(&stream_id) else {
return StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
);
};
if self.expire_stream_if_due(&stream_id, now_ms) {
return StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
);
}
let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
Ok(decision) => decision,
Err(response) => return response,
};
if let ProducerDecision::Duplicate {
offset,
next_offset,
closed,
producer,
..
} = producer_decision
{
if payload.is_empty() {
return StreamResponse::Closed {
next_offset,
deduplicated: true,
producer: Some(producer),
};
}
return StreamResponse::Appended {
offset,
next_offset,
closed,
deduplicated: true,
producer: Some(producer),
};
}
if let Err(response) = self.validate_record_match(&stream_id, record_match) {
return response;
}
let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
let record_ends = match content_type {
Some(value) => match canonical_json_record_ends(value, payload) {
Ok(record_ends) => record_ends,
Err(_) => {
return StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"application/json append payload must use canonical newline boundaries",
);
}
},
None => Vec::new(),
};
let prepared_record_append = {
let slot = self
.stream_slot(&stream_id)
.expect("stream existence checked before record validation");
match prepare_record_append(
slot.record_index.as_ref(),
super::is_json_record_content_type(&slot.metadata.content_type),
slot.metadata.tail_offset,
payload_len,
&record_ends,
) {
Ok(prepared) => prepared,
Err(response) => return response,
}
};
let record_range = prepared_record_append
.as_ref()
.map(crate::PreparedRecordAppend::range);
if let Err(response) = self.check_append_quota(&stream_id.bucket_id, payload_len) {
return response;
}
let Some(stream) = self.stream_metadata_mut(&stream_id) else {
unreachable!("stream existence checked before producer evaluation");
};
if stream.status == StreamStatus::Closed {
if close_after && payload.is_empty() {
return StreamResponse::Closed {
next_offset: stream.tail_offset,
deduplicated: false,
producer: None,
};
}
return StreamResponse::error_with_next_offset_and_context(
StreamErrorCode::StreamClosed,
format!("stream '{stream_id}' is closed"),
stream.tail_offset,
vec![StreamErrorContext::StreamClosed],
);
}
if payload.is_empty() && !close_after {
return StreamResponse::error(
StreamErrorCode::EmptyAppend,
"append payload must be non-empty unless closing the stream",
);
}
if !payload.is_empty() {
let Some(content_type) = content_type else {
return StreamResponse::error(
StreamErrorCode::MissingContentType,
"append with a body must include content type",
);
};
if content_type != stream.content_type {
return StreamResponse::error_with_next_offset(
StreamErrorCode::ContentTypeMismatch,
format!(
"append content type '{content_type}' does not match stream content type '{}'",
stream.content_type
),
stream.tail_offset,
);
}
}
if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
return response;
}
let offset = stream.tail_offset;
stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
if let Some(seq) = stream_seq {
stream.last_stream_seq = Some(seq);
}
renew_stream_ttl(stream, now_ms);
if close_after {
stream.status = StreamStatus::Closed;
}
let closed = stream.status == StreamStatus::Closed;
let next_offset = stream.tail_offset;
self.refresh_ttl_entry(&stream_id);
let producer_ack = producer.clone();
if let Some(producer) = producer {
self.record_producer_success(
stream_id.clone(),
producer,
ProducerAppendRecord {
start_offset: offset,
next_offset,
closed,
record_start: record_range.map(|range| range.first_record),
record_next: record_range.map(|range| range.next_record),
},
vec![ProducerAppendRecord {
start_offset: offset,
next_offset,
closed,
record_start: record_range.map(|range| range.first_record),
record_next: record_range.map(|range| range.next_record),
}],
);
}
if payload.is_empty() {
StreamResponse::Closed {
next_offset,
deduplicated: false,
producer: producer_ack,
}
} else {
let slot = self
.stream_slot_mut(&stream_id)
.expect("stream existence checked before append mutation");
if let (Some(index), Some(prepared)) =
(slot.record_index.as_mut(), prepared_record_append)
{
let _range = index.commit_append(prepared);
}
slot.hot_buffer.push(offset, next_offset, payload);
slot.integrity
.append_payload(&stream_id, offset, next_offset, payload);
slot.message_records
.extend(Self::message_records_for_append(
offset,
next_offset,
&record_ends,
));
self.add_hot_payload_bytes(payload_len);
self.usage_on_append(
&stream_id.bucket_id,
payload_len,
Self::appended_record_count(&record_ends, payload_len),
);
StreamResponse::Appended {
offset,
next_offset,
closed: close_after,
deduplicated: false,
producer: producer_ack,
}
}
}
pub(super) fn append_external(&mut self, input: AppendExternalInput<'_>) -> StreamResponse {
let AppendExternalInput {
stream_id,
content_type,
payload,
record_ends,
close_after,
stream_seq,
producer,
now_ms,
record_match,
} = input;
if let Err(response) = validate_external_payload_ref(&payload) {
return response;
}
if let Err(response) = self.validate_stream_scope(&stream_id) {
return response;
}
if let Err(response) = validate_producer_request(producer.as_ref()) {
return response;
}
let Some(_) = self.stream_metadata(&stream_id) else {
return StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
);
};
if self.expire_stream_if_due(&stream_id, now_ms) {
return StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
);
}
let producer_decision = match self.evaluate_producer(&stream_id, producer.as_ref()) {
Ok(decision) => decision,
Err(response) => return response,
};
if let ProducerDecision::Duplicate {
offset,
next_offset,
closed,
producer,
..
} = producer_decision
{
return StreamResponse::Appended {
offset,
next_offset,
closed,
deduplicated: true,
producer: Some(producer),
};
}
if let Err(response) = self.validate_record_match(&stream_id, record_match) {
return response;
}
let prepared_record_append = {
let slot = self
.stream_slot(&stream_id)
.expect("stream existence checked before record validation");
match prepare_record_append(
slot.record_index.as_ref(),
super::is_json_record_content_type(&slot.metadata.content_type),
slot.metadata.tail_offset,
payload.payload_len,
&record_ends,
) {
Ok(prepared) => prepared,
Err(response) => return response,
}
};
let record_range = prepared_record_append
.as_ref()
.map(crate::PreparedRecordAppend::range);
let Some(stream) = self.stream_metadata(&stream_id) else {
unreachable!("stream existence checked before producer evaluation");
};
if stream.status == StreamStatus::Closed {
return StreamResponse::error_with_next_offset_and_context(
StreamErrorCode::StreamClosed,
format!("stream '{stream_id}' is closed"),
stream.tail_offset,
vec![StreamErrorContext::StreamClosed],
);
}
let Some(content_type) = content_type else {
return StreamResponse::error(
StreamErrorCode::MissingContentType,
"append with a body must include content type",
);
};
if content_type != stream.content_type {
return StreamResponse::error_with_next_offset(
StreamErrorCode::ContentTypeMismatch,
format!(
"append content type '{content_type}' does not match stream content type '{}'",
stream.content_type
),
stream.tail_offset,
);
}
if let Err(response) = check_stream_seq(stream, stream_seq.as_deref()) {
return response;
}
let offset = stream.tail_offset;
let next_offset = offset.saturating_add(payload.payload_len);
if let Err(response) = self.check_append_quota(&stream_id.bucket_id, payload.payload_len) {
return response;
}
let stream = self
.stream_metadata_mut(&stream_id)
.expect("stream existence checked before external append mutation");
stream.tail_offset = next_offset;
if let Some(seq) = stream_seq {
stream.last_stream_seq = Some(seq);
}
renew_stream_ttl(stream, now_ms);
if close_after {
stream.status = StreamStatus::Closed;
}
let closed = stream.status == StreamStatus::Closed;
self.refresh_ttl_entry(&stream_id);
let producer_ack = producer.clone();
if let Some(producer) = producer {
self.record_producer_success(
stream_id.clone(),
producer,
ProducerAppendRecord {
start_offset: offset,
next_offset,
closed,
record_start: record_range.map(|range| range.first_record),
record_next: record_range.map(|range| range.next_record),
},
vec![ProducerAppendRecord {
start_offset: offset,
next_offset,
closed,
record_start: record_range.map(|range| range.first_record),
record_next: record_range.map(|range| range.next_record),
}],
);
}
let object = ObjectPayloadRef {
start_offset: offset,
end_offset: next_offset,
s3_path: payload.s3_path,
object_size: payload.object_size,
object_offset: 0,
};
let slot = self
.stream_slot_mut(&stream_id)
.expect("stream existence checked before external append mutation");
if let (Some(index), Some(prepared)) = (slot.record_index.as_mut(), prepared_record_append)
{
let _range = index.commit_append(prepared);
}
slot.cold.push_external_segment(object.clone());
slot.integrity.append_external(
&stream_id,
object.start_offset,
object.end_offset,
&object.s3_path,
object.object_size,
);
slot.message_records
.extend(Self::message_records_for_append(
offset,
next_offset,
&record_ends,
));
let appended_bytes = next_offset.saturating_sub(offset);
self.usage_on_append(
&stream_id.bucket_id,
appended_bytes,
Self::appended_record_count(&record_ends, appended_bytes),
);
StreamResponse::Appended {
offset,
next_offset,
closed: close_after,
deduplicated: false,
producer: producer_ack,
}
}
pub fn append_batch_borrowed(
&mut self,
stream_id: BucketStreamId,
content_type: Option<&str>,
payloads: &[&[u8]],
producer: Option<ProducerRequest>,
now_ms: u64,
) -> Result<StreamBatchAppend, StreamResponse> {
if payloads.is_empty() {
return Err(StreamResponse::error(
StreamErrorCode::EmptyAppend,
"append batch must contain at least one payload",
));
}
self.validate_stream_scope(&stream_id)?;
validate_producer_request(producer.as_ref())?;
if self.expire_stream_if_due(&stream_id, now_ms) {
return Err(StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
));
}
let producer_decision = self.evaluate_producer(&stream_id, producer.as_ref())?;
if let ProducerDecision::Duplicate { items, .. } = producer_decision {
return Ok(StreamBatchAppend {
items: items
.into_iter()
.map(|item| StreamBatchAppendItem {
offset: item.start_offset,
next_offset: item.next_offset,
closed: item.closed,
deduplicated: true,
})
.collect(),
deduplicated: true,
});
}
let Some(stream) = self.stream_metadata(&stream_id) else {
return Err(StreamResponse::error(
StreamErrorCode::StreamNotFound,
format!("stream '{stream_id}' does not exist"),
));
};
if stream.status == StreamStatus::Closed {
return Err(StreamResponse::error_with_next_offset_and_context(
StreamErrorCode::StreamClosed,
format!("stream '{stream_id}' is closed"),
stream.tail_offset,
vec![StreamErrorContext::StreamClosed],
));
}
let Some(content_type) = content_type else {
return Err(StreamResponse::error(
StreamErrorCode::MissingContentType,
"append batch must include content type",
));
};
if content_type != stream.content_type {
return Err(StreamResponse::error_with_next_offset(
StreamErrorCode::ContentTypeMismatch,
format!(
"append content type '{content_type}' does not match stream content type '{}'",
stream.content_type
),
stream.tail_offset,
));
}
if payloads.iter().any(|payload| payload.is_empty()) {
return Err(StreamResponse::error(
StreamErrorCode::EmptyAppend,
"append batch payloads must be non-empty",
));
}
let base_offset = stream.tail_offset;
let mut total_payload_len = 0_u64;
let mut combined_record_ends = Vec::new();
let mut record_counts = Vec::with_capacity(payloads.len());
let mut all_record_ends = Vec::with_capacity(payloads.len());
for payload in payloads {
let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
let record_ends = canonical_json_record_ends(content_type, payload).map_err(|_| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"application/json append payload must use canonical newline boundaries",
)
})?;
record_counts.push(u64::try_from(record_ends.len()).map_err(|_| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"record count exceeds the supported range",
)
})?);
for end in &record_ends {
combined_record_ends.push(total_payload_len.checked_add(*end).ok_or_else(
|| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"append batch payload length exceeds the supported range",
)
},
)?);
}
total_payload_len = total_payload_len.checked_add(payload_len).ok_or_else(|| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"append batch payload length exceeds the supported range",
)
})?;
all_record_ends.push(record_ends);
}
let prepared_record_append = {
let slot = self
.stream_slot(&stream_id)
.expect("stream existence checked before record validation");
prepare_record_append(
slot.record_index.as_ref(),
super::is_json_record_content_type(content_type),
base_offset,
total_payload_len,
&combined_record_ends,
)?
};
let mut next_record = prepared_record_append
.as_ref()
.map(crate::PreparedRecordAppend::range)
.map(|range| range.first_record);
let mut item_record_ranges = Vec::with_capacity(record_counts.len());
for count in record_counts {
let range = match next_record {
Some(first_record) => {
let next = first_record.checked_add(count).ok_or_else(|| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"append batch record count exceeds the supported range",
)
})?;
next_record = Some(next);
Some(crate::StreamRecordRange {
first_record,
next_record: next,
})
}
None => None,
};
item_record_ranges.push(range);
}
let batch_bytes = payloads.iter().fold(0u64, |total, payload| {
total.saturating_add(payload.len() as u64)
});
self.check_append_quota(&stream_id.bucket_id, batch_bytes)?;
let stream = self
.stream_metadata_mut(&stream_id)
.expect("stream existence checked before batch append mutation");
let mut items = Vec::with_capacity(payloads.len());
for (payload, record_range) in payloads.iter().zip(item_record_ranges) {
let offset = stream.tail_offset;
let payload_len = u64::try_from(payload.len()).expect("payload len fits u64");
stream.tail_offset = stream.tail_offset.saturating_add(payload_len);
items.push(ProducerAppendRecord {
start_offset: offset,
next_offset: stream.tail_offset,
closed: false,
record_start: record_range.map(|range| range.first_record),
record_next: record_range.map(|range| range.next_record),
});
}
let last = items
.last()
.expect("payloads checked non-empty before append")
.clone();
renew_stream_ttl(stream, now_ms);
self.refresh_ttl_entry(&stream_id);
if let Some(producer) = producer {
self.record_producer_success(stream_id.clone(), producer, last.clone(), items.clone());
}
let slot = self
.stream_slot_mut(&stream_id)
.expect("stream existence checked before batch append mutation");
if let (Some(index), Some(prepared)) = (slot.record_index.as_mut(), prepared_record_append)
{
let _range = index.commit_append(prepared);
}
for (item, payload) in items.iter().zip(payloads.iter()) {
slot.hot_buffer
.push(item.start_offset, item.next_offset, payload);
}
for (item, payload) in items.iter().zip(payloads.iter()) {
slot.integrity
.append_payload(&stream_id, item.start_offset, item.next_offset, payload);
}
for (item, record_ends) in items.iter().zip(all_record_ends.iter()) {
slot.message_records
.extend(Self::message_records_for_append(
item.start_offset,
item.next_offset,
record_ends,
));
}
let mut appended_bytes: u64 = 0;
let mut appended_records: u64 = 0;
for (item, record_ends) in items.iter().zip(all_record_ends.iter()) {
let item_bytes = item.next_offset.saturating_sub(item.start_offset);
appended_bytes = appended_bytes.saturating_add(item_bytes);
appended_records = appended_records
.saturating_add(Self::appended_record_count(record_ends, item_bytes));
}
self.add_hot_payload_bytes(appended_bytes);
self.usage_on_append(&stream_id.bucket_id, appended_bytes, appended_records);
Ok(StreamBatchAppend {
items: items
.into_iter()
.map(|item| StreamBatchAppendItem {
offset: item.start_offset,
next_offset: item.next_offset,
closed: item.closed,
deduplicated: false,
})
.collect(),
deduplicated: false,
})
}
fn validate_record_match(
&self,
stream_id: &BucketStreamId,
expected: Option<u64>,
) -> Result<(), StreamResponse> {
let Some(expected) = expected else {
return Ok(());
};
let Some(slot) = self.stream_slot(stream_id) else {
return Ok(());
};
let Some(index) = slot.record_index.as_ref() else {
return Err(StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"Stream-Record-Match requires active JSON record coordinates",
));
};
let current = index
.range()
.map_err(|_| {
StreamResponse::error(
StreamErrorCode::InvalidRecordBoundaries,
"stream record index is invalid",
)
})?
.next_record;
if current == expected {
return Ok(());
}
Err(StreamResponse::error_with_next_offset_and_context(
StreamErrorCode::RecordPreconditionFailed,
format!("record tail is {current}, expected {expected}"),
slot.metadata.tail_offset,
vec![StreamErrorContext::RecordTailMismatch {
current_record: current,
}],
))
}
fn evaluate_producer(
&self,
stream_id: &BucketStreamId,
producer: Option<&ProducerRequest>,
) -> Result<ProducerDecision, StreamResponse> {
let Some(producer) = producer else {
return Ok(ProducerDecision::Accept);
};
let Some(states) = self.stream_slot(stream_id).map(|slot| &slot.producers) else {
return Ok(ProducerDecision::Accept);
};
let Some(state) = states.get(&producer.producer_id) else {
if producer.producer_seq == 0 {
return Ok(ProducerDecision::Accept);
}
return Err(StreamResponse::error_with_context(
StreamErrorCode::ProducerSeqConflict,
format!(
"producer '{}' expected sequence 0, received {}",
producer.producer_id, producer.producer_seq
),
vec![StreamErrorContext::ProducerSeqConflict {
expected_seq: 0,
received_seq: producer.producer_seq,
}],
));
};
if producer.producer_epoch < state.producer_epoch {
return Err(StreamResponse::error_with_context(
StreamErrorCode::ProducerEpochStale,
format!(
"producer '{}' epoch {} is stale; current epoch is {}",
producer.producer_id, producer.producer_epoch, state.producer_epoch
),
vec![StreamErrorContext::ProducerEpochStale {
current_epoch: state.producer_epoch,
}],
));
}
if producer.producer_epoch > state.producer_epoch {
if producer.producer_seq == 0 {
return Ok(ProducerDecision::Accept);
}
return Err(StreamResponse::error(
StreamErrorCode::InvalidProducer,
format!(
"producer '{}' new epoch {} must start at sequence 0",
producer.producer_id, producer.producer_epoch
),
));
}
if producer.producer_seq <= state.producer_seq {
let Some(receipt) = state
.receipts
.iter()
.find(|receipt| receipt.producer_seq == producer.producer_seq)
else {
return Err(StreamResponse::error_with_context(
StreamErrorCode::ProducerSeqConflict,
format!(
"producer '{}' sequence {} is older than the retained receipt window ending at {}",
producer.producer_id, producer.producer_seq, state.producer_seq
),
vec![StreamErrorContext::ProducerSeqConflict {
expected_seq: state.producer_seq.saturating_add(1),
received_seq: producer.producer_seq,
}],
));
};
return Ok(ProducerDecision::Duplicate {
offset: receipt.start_offset,
next_offset: receipt.next_offset,
closed: receipt.closed,
producer: ProducerRequest {
producer_id: producer.producer_id.clone(),
producer_epoch: state.producer_epoch,
producer_seq: receipt.producer_seq,
},
items: receipt.items.clone(),
});
}
if producer.producer_seq == state.producer_seq + 1 {
return Ok(ProducerDecision::Accept);
}
Err(StreamResponse::error_with_context(
StreamErrorCode::ProducerSeqConflict,
format!(
"producer '{}' expected sequence {}, received {}",
producer.producer_id,
state.producer_seq + 1,
producer.producer_seq
),
vec![StreamErrorContext::ProducerSeqConflict {
expected_seq: state.producer_seq + 1,
received_seq: producer.producer_seq,
}],
))
}
fn record_producer_success(
&mut self,
stream_id: BucketStreamId,
producer: ProducerRequest,
last: ProducerAppendRecord,
last_items: Vec<ProducerAppendRecord>,
) {
let receipt = ProducerReceipt {
producer_seq: producer.producer_seq,
start_offset: last.start_offset,
next_offset: last.next_offset,
closed: last.closed,
items: last_items.clone(),
};
let producers = &mut self
.stream_slot_mut(&stream_id)
.expect("stream existence checked before producer mutation")
.producers;
if let Some(state) = producers.get_mut(&producer.producer_id)
&& state.producer_epoch == producer.producer_epoch
{
state.producer_seq = producer.producer_seq;
state.last_start_offset = last.start_offset;
state.last_next_offset = last.next_offset;
state.last_closed = last.closed;
state.last_items = last_items;
state.receipts.push(receipt);
return;
}
producers.insert(producer.producer_id, ProducerState {
producer_epoch: producer.producer_epoch,
producer_seq: producer.producer_seq,
last_start_offset: last.start_offset,
last_next_offset: last.next_offset,
last_closed: last.closed,
last_items,
receipts: vec![receipt],
});
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum ProducerDecision {
Accept,
Duplicate {
offset: u64,
next_offset: u64,
closed: bool,
producer: ProducerRequest,
items: Vec<ProducerAppendRecord>,
},
}
fn check_stream_seq(stream: &StreamMetadata, incoming: Option<&str>) -> Result<(), StreamResponse> {
let Some(incoming) = incoming else {
return Ok(());
};
if let Some(last) = stream.last_stream_seq.as_deref()
&& incoming <= last
{
return Err(StreamResponse::error_with_next_offset(
StreamErrorCode::StreamSeqConflict,
format!("stream sequence '{incoming}' is not greater than last sequence '{last}'"),
stream.tail_offset,
));
}
Ok(())
}