use serde::{Deserialize, Serialize};
use crate::message::ContentPart;
pub const MAX_SUBMISSION_ITEM_BYTES: usize = 64 * 1024;
pub const MAX_STEERING_QUEUE_ITEMS: usize = 128;
pub const MAX_STEERING_QUEUE_BYTES: usize = 1024 * 1024;
pub const MAX_STEERING_QUEUE_IMAGES: usize = 16;
pub const MAX_STEERING_QUEUE_IMAGE_BYTES: u64 = 100 * 1024 * 1024;
pub const MAX_SUBMISSION_BATCH_ITEMS: usize = 32;
pub const MAX_SUBMISSION_BATCH_BYTES: usize = 256 * 1024;
pub const MAX_SUBMISSION_ATTACHMENT_METADATA_BYTES: usize = 16 * 1024;
pub const MAX_SUBMISSION_TOTAL_ATTACHMENT_METADATA_BYTES: usize = 64 * 1024;
pub const MAX_SUBMISSION_IMAGE_COUNT: usize = 4;
pub const MAX_SUBMISSION_IMAGE_BYTES: u64 = 20 * 1024 * 1024;
pub const MAX_SUBMISSION_TOTAL_IMAGE_BYTES: u64 = 50 * 1024 * 1024;
pub const MAX_PENDING_SUBMISSIONS: usize = 128;
pub const MAX_PENDING_SUBMISSION_BYTES: usize = 1024 * 1024;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SubmissionSource {
User,
Scheduler,
Compatibility,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SubmissionKind {
UserTurn,
PreviewRequest,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct SubmissionItem {
pub id: String,
pub enqueue_sequence: u64,
pub kind: SubmissionKind,
pub text: String,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub attachments: Vec<ContentPart>,
}
impl SubmissionItem {
#[must_use]
pub fn text_bytes(&self) -> usize {
self.text.len()
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct StructuredSubmission {
pub id: String,
pub source: SubmissionSource,
#[serde(default)]
pub sender_generation: u64,
pub items: Vec<SubmissionItem>,
}
impl StructuredSubmission {
#[must_use]
pub fn total_text_bytes(&self) -> usize {
self.items.iter().fold(0usize, |total, item| {
total.saturating_add(item.text_bytes())
})
}
#[must_use]
pub fn image_totals(&self) -> (usize, u64) {
self.items.iter().flat_map(|item| &item.attachments).fold(
(0usize, 0u64),
|(count, bytes), part| match part {
ContentPart::Image { byte_count, .. } => {
(count.saturating_add(1), bytes.saturating_add(*byte_count))
}
ContentPart::Text { .. } => (count, bytes),
},
)
}
#[must_use]
pub fn attachment_metadata_bytes(&self) -> Option<usize> {
self.items
.iter()
.flat_map(|item| &item.attachments)
.try_fold(0usize, |total, part| match part {
ContentPart::Image { path, mime, .. } => total
.checked_add(path.as_os_str().to_string_lossy().len())?
.checked_add(mime.len())?
.checked_add(std::mem::size_of::<u64>())?
.checked_add(32),
ContentPart::Text { .. } => None,
})
}
#[must_use]
pub fn common_kind(&self) -> Option<SubmissionKind> {
let first = self.items.first()?.kind;
self.items
.iter()
.all(|item| item.kind == first)
.then_some(first)
}
pub fn validate(&self) -> Result<(), SubmissionRejectionReason> {
if self.id.is_empty()
|| self.items.is_empty()
|| self.items.len() > MAX_SUBMISSION_BATCH_ITEMS
|| self.common_kind().is_none()
{
return Err(SubmissionRejectionReason::InvalidStructure);
}
let mut previous_sequence = None;
let mut item_ids = std::collections::HashSet::with_capacity(self.items.len());
for item in &self.items {
if item.id.is_empty() {
return Err(SubmissionRejectionReason::InvalidStructure);
}
if !item_ids.insert(item.id.as_str()) {
return Err(SubmissionRejectionReason::Duplicate);
}
if item.text_bytes() > MAX_SUBMISSION_ITEM_BYTES {
return Err(SubmissionRejectionReason::LimitExceeded);
}
if previous_sequence.is_some_and(|previous| item.enqueue_sequence <= previous) {
return Err(SubmissionRejectionReason::InvalidStructure);
}
previous_sequence = Some(item.enqueue_sequence);
}
if self.total_text_bytes() > MAX_SUBMISSION_BATCH_BYTES {
return Err(SubmissionRejectionReason::LimitExceeded);
}
if self.common_kind() == Some(SubmissionKind::PreviewRequest)
&& (self.items.len() != 1 || !self.items[0].attachments.is_empty())
{
return Err(SubmissionRejectionReason::InvalidStructure);
}
let (image_count, image_bytes) = self.image_totals();
let metadata_bytes = self
.attachment_metadata_bytes()
.ok_or(SubmissionRejectionReason::InvalidStructure)?;
let oversized_attachment =
self.items
.iter()
.flat_map(|item| &item.attachments)
.any(|part| match part {
ContentPart::Image {
path,
mime,
byte_count,
..
} => {
let metadata = path
.as_os_str()
.to_string_lossy()
.len()
.saturating_add(mime.len())
.saturating_add(std::mem::size_of::<u64>())
.saturating_add(32);
*byte_count > MAX_SUBMISSION_IMAGE_BYTES
|| metadata > MAX_SUBMISSION_ATTACHMENT_METADATA_BYTES
}
ContentPart::Text { .. } => true,
});
if image_count > MAX_SUBMISSION_IMAGE_COUNT
|| image_bytes > MAX_SUBMISSION_TOTAL_IMAGE_BYTES
|| metadata_bytes > MAX_SUBMISSION_TOTAL_ATTACHMENT_METADATA_BYTES
|| oversized_attachment
{
return Err(SubmissionRejectionReason::LimitExceeded);
}
Ok(())
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum PendingSubmissionState {
AcceptedPending,
Running,
PausedPending,
TerminalCancelled,
TerminalError,
Committed,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SubmissionRejectionReason {
DurabilityUnavailable,
WrongSession,
WrongGeneration,
IdentityConflict,
Duplicate,
Cancelled,
InvalidStructure,
LimitExceeded,
ContextBudgetExceeded,
SessionClosed,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(tag = "status", rename_all = "snake_case")]
pub enum SubmissionReceiptDisposition {
AcceptedPending,
AlreadyAccepted {
state: PendingSubmissionState,
#[serde(default, skip_serializing_if = "Option::is_none")]
turn_id: Option<String>,
},
NotAccepted,
Rejected {
reason: SubmissionRejectionReason,
},
}
impl SubmissionReceiptDisposition {
#[must_use]
pub fn has_durable_custody(&self) -> bool {
matches!(self, Self::AcceptedPending | Self::AlreadyAccepted { .. })
}
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct SubmissionReceipt {
pub session_id: String,
pub session_generation: u64,
pub submission_id: String,
pub reservation_id: String,
pub receipt_id: String,
pub source: SubmissionSource,
pub item_count: usize,
pub total_text_bytes: usize,
pub disposition: SubmissionReceiptDisposition,
}
#[cfg(test)]
mod tests {
use super::*;
fn submission() -> StructuredSubmission {
StructuredSubmission {
id: "batch-1".into(),
source: SubmissionSource::User,
sender_generation: 1,
items: vec![SubmissionItem {
id: "item-1".into(),
enqueue_sequence: 1,
kind: SubmissionKind::UserTurn,
text: "hello".into(),
attachments: Vec::new(),
}],
}
}
#[test]
fn valid_submission_preserves_item_boundaries() {
let mut value = submission();
value.items.push(SubmissionItem {
id: "item-2".into(),
enqueue_sequence: 2,
kind: SubmissionKind::UserTurn,
text: "world".into(),
attachments: Vec::new(),
});
assert_eq!(value.validate(), Ok(()));
let encoded = serde_json::to_string(&value).expect("operation should succeed");
let decoded: StructuredSubmission =
serde_json::from_str(&encoded).expect("operation should succeed");
assert_eq!(decoded.items.len(), 2);
assert_eq!(decoded.items[0].text, "hello");
assert_eq!(decoded.items[1].text, "world");
}
#[test]
fn duplicate_item_identity_fails_closed() {
let mut value = submission();
value.items.push(SubmissionItem {
id: "item-1".into(),
enqueue_sequence: 2,
kind: SubmissionKind::UserTurn,
text: "again".into(),
attachments: Vec::new(),
});
assert_eq!(value.validate(), Err(SubmissionRejectionReason::Duplicate));
}
#[test]
fn incompatible_kinds_and_regressive_sequence_fail_closed() {
let mut value = submission();
value.items.push(SubmissionItem {
id: "item-2".into(),
enqueue_sequence: 1,
kind: SubmissionKind::PreviewRequest,
text: "preview".into(),
attachments: Vec::new(),
});
assert_eq!(
value.validate(),
Err(SubmissionRejectionReason::InvalidStructure)
);
}
#[test]
fn durable_custody_excludes_rejection_and_not_accepted() {
assert!(SubmissionReceiptDisposition::AcceptedPending.has_durable_custody());
assert!(
SubmissionReceiptDisposition::AlreadyAccepted {
state: PendingSubmissionState::Committed,
turn_id: Some("turn-1".into()),
}
.has_durable_custody()
);
assert!(!SubmissionReceiptDisposition::NotAccepted.has_durable_custody());
assert!(
!SubmissionReceiptDisposition::Rejected {
reason: SubmissionRejectionReason::LimitExceeded,
}
.has_durable_custody()
);
}
}