use serde::{Deserialize, Serialize};
use sha2::{Digest, Sha256};
use std::collections::BTreeMap;
use std::path::{Path, PathBuf};
use tokio::io::{AsyncRead, AsyncReadExt, AsyncSeekExt, AsyncWrite, AsyncWriteExt};
const ACK_OK: u8 = 0x01;
const ACK_REJECTED: u8 = 0x02;
pub const DEFAULT_CHUNK_BYTES: usize = 64 * 1024;
pub const DEFAULT_MAX_HEADER_BYTES: usize = 64 * 1024;
pub const DEFAULT_MAX_FILE_BYTES: u64 = 256 * 1024 * 1024;
pub const ROUTED_FRAME_MAGIC: &[u8; 4] = b"OFT2";
pub const DEFAULT_ROUTED_CHUNK_BYTES: usize = 48 * 1024;
pub const MAX_ROUTED_TRANSFER_ID_BYTES: usize = 512;
const ROUTED_KIND_OFFER: u8 = 0x01;
const ROUTED_KIND_RESUME: u8 = 0x02;
const ROUTED_KIND_CHUNK: u8 = 0x03;
const ROUTED_KIND_ACK: u8 = 0x04;
const ROUTED_KIND_COMPLETE: u8 = 0x05;
const ROUTED_KIND_COMPLETED: u8 = 0x06;
const ROUTED_KIND_REJECT: u8 = 0x07;
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum RoutedTransferFrame {
Offer(FileTransferHeader),
Resume {
transfer_id: String,
offset: u64,
},
Chunk {
transfer_id: String,
offset: u64,
ack_requested: bool,
bytes: Vec<u8>,
},
Ack {
transfer_id: String,
offset: u64,
},
Complete {
transfer_id: String,
size: u64,
},
Completed {
transfer_id: String,
size: u64,
},
Reject {
transfer_id: String,
reason: String,
},
}
impl RoutedTransferFrame {
pub fn transfer_id(&self) -> &str {
match self {
Self::Offer(header) => &header.transfer_id,
Self::Resume { transfer_id, .. }
| Self::Chunk { transfer_id, .. }
| Self::Ack { transfer_id, .. }
| Self::Complete { transfer_id, .. }
| Self::Completed { transfer_id, .. }
| Self::Reject { transfer_id, .. } => transfer_id,
}
}
pub fn encode(&self, limits: TransferLimits) -> Result<Vec<u8>, TransferError> {
let (kind, transfer_id) = match self {
Self::Offer(header) => (ROUTED_KIND_OFFER, header.transfer_id.as_str()),
Self::Resume { transfer_id, .. } => (ROUTED_KIND_RESUME, transfer_id.as_str()),
Self::Chunk { transfer_id, .. } => (ROUTED_KIND_CHUNK, transfer_id.as_str()),
Self::Ack { transfer_id, .. } => (ROUTED_KIND_ACK, transfer_id.as_str()),
Self::Complete { transfer_id, .. } => (ROUTED_KIND_COMPLETE, transfer_id.as_str()),
Self::Completed { transfer_id, .. } => (ROUTED_KIND_COMPLETED, transfer_id.as_str()),
Self::Reject { transfer_id, .. } => (ROUTED_KIND_REJECT, transfer_id.as_str()),
};
let id = transfer_id.as_bytes();
if id.is_empty() || id.len() > MAX_ROUTED_TRANSFER_ID_BYTES || id.len() > u16::MAX as usize
{
return Err(TransferError::Limit(
"routed transferId length is invalid".into(),
));
}
let mut output = Vec::with_capacity(7 + id.len() + 16);
output.extend_from_slice(ROUTED_FRAME_MAGIC);
output.push(kind);
output.extend_from_slice(&(id.len() as u16).to_be_bytes());
output.extend_from_slice(id);
match self {
Self::Offer(header) => {
header.validate(limits)?;
let encoded = serde_json::to_vec(header)?;
if encoded.is_empty()
|| encoded.len() > limits.max_header_bytes
|| encoded.len() > u32::MAX as usize
{
return Err(TransferError::Limit(
"routed offer header is too large".into(),
));
}
output.extend_from_slice(&(encoded.len() as u32).to_be_bytes());
output.extend_from_slice(&encoded);
}
Self::Resume { offset, .. } | Self::Ack { offset, .. } => {
output.extend_from_slice(&offset.to_be_bytes());
}
Self::Chunk {
offset,
ack_requested,
bytes,
..
} => {
if bytes.is_empty() || bytes.len() > limits.chunk_bytes {
return Err(TransferError::Limit(format!(
"routed chunk must contain 1..={} bytes",
limits.chunk_bytes
)));
}
output.extend_from_slice(&offset.to_be_bytes());
output.push(u8::from(*ack_requested));
output.extend_from_slice(bytes);
}
Self::Complete { size, .. } | Self::Completed { size, .. } => {
output.extend_from_slice(&size.to_be_bytes());
}
Self::Reject { reason, .. } => {
let bytes = reason.as_bytes();
if bytes.len() > limits.max_header_bytes || bytes.len() > u16::MAX as usize {
return Err(TransferError::Limit(
"routed rejection reason is too large".into(),
));
}
output.extend_from_slice(&(bytes.len() as u16).to_be_bytes());
output.extend_from_slice(bytes);
}
}
Ok(output)
}
pub fn decode(input: &[u8], limits: TransferLimits) -> Result<Self, TransferError> {
if input.len() < 7 || &input[..4] != ROUTED_FRAME_MAGIC {
return Err(TransferError::Protocol(
"not a routed file-transfer frame".into(),
));
}
let kind = input[4];
let id_len = u16::from_be_bytes([input[5], input[6]]) as usize;
if id_len == 0 || id_len > MAX_ROUTED_TRANSFER_ID_BYTES || input.len() < 7 + id_len {
return Err(TransferError::Protocol("invalid routed transferId".into()));
}
let transfer_id = std::str::from_utf8(&input[7..7 + id_len])
.map_err(|_| TransferError::Protocol("routed transferId is not UTF-8".into()))?
.to_string();
let body = &input[7 + id_len..];
let read_u64 = |bytes: &[u8]| -> Result<u64, TransferError> {
let bytes: [u8; 8] = bytes
.get(..8)
.ok_or_else(|| TransferError::Protocol("truncated routed frame".into()))?
.try_into()
.map_err(|_| TransferError::Protocol("truncated routed frame".into()))?;
Ok(u64::from_be_bytes(bytes))
};
match kind {
ROUTED_KIND_OFFER => {
let len_bytes: [u8; 4] = body
.get(..4)
.ok_or_else(|| TransferError::Protocol("truncated routed offer".into()))?
.try_into()
.map_err(|_| TransferError::Protocol("truncated routed offer".into()))?;
let len = u32::from_be_bytes(len_bytes) as usize;
if len == 0 || len > limits.max_header_bytes || body.len() != 4 + len {
return Err(TransferError::Limit(
"invalid routed offer header length".into(),
));
}
let header: FileTransferHeader = serde_json::from_slice(&body[4..])?;
header.validate(limits)?;
if header.transfer_id != transfer_id {
return Err(TransferError::Protocol(
"routed offer transferId mismatch".into(),
));
}
Ok(Self::Offer(header))
}
ROUTED_KIND_RESUME if body.len() == 8 => Ok(Self::Resume {
transfer_id,
offset: read_u64(body)?,
}),
ROUTED_KIND_ACK if body.len() == 8 => Ok(Self::Ack {
transfer_id,
offset: read_u64(body)?,
}),
ROUTED_KIND_COMPLETE if body.len() == 8 => Ok(Self::Complete {
transfer_id,
size: read_u64(body)?,
}),
ROUTED_KIND_COMPLETED if body.len() == 8 => Ok(Self::Completed {
transfer_id,
size: read_u64(body)?,
}),
ROUTED_KIND_CHUNK if body.len() > 9 => {
let offset = read_u64(body)?;
let ack_requested = match body[8] {
0 => false,
1 => true,
_ => {
return Err(TransferError::Protocol(
"invalid routed chunk ACK flag".into(),
))
}
};
let bytes = body[9..].to_vec();
if bytes.len() > limits.chunk_bytes {
return Err(TransferError::Limit(
"routed chunk exceeds configured limit".into(),
));
}
Ok(Self::Chunk {
transfer_id,
offset,
ack_requested,
bytes,
})
}
ROUTED_KIND_REJECT => {
let len_bytes: [u8; 2] = body
.get(..2)
.ok_or_else(|| TransferError::Protocol("truncated routed rejection".into()))?
.try_into()
.map_err(|_| TransferError::Protocol("truncated routed rejection".into()))?;
let len = u16::from_be_bytes(len_bytes) as usize;
if body.len() != 2 + len {
return Err(TransferError::Protocol(
"invalid routed rejection length".into(),
));
}
let reason = std::str::from_utf8(&body[2..])
.map_err(|_| TransferError::Protocol("routed rejection is not UTF-8".into()))?
.to_string();
Ok(Self::Reject {
transfer_id,
reason,
})
}
_ => Err(TransferError::Protocol(
"invalid routed file-transfer frame".into(),
)),
}
}
}
#[derive(Debug, Clone, Serialize, Deserialize, PartialEq, Eq)]
#[serde(rename_all = "camelCase")]
pub struct FileTransferHeader {
pub version: u8,
pub transfer_id: String,
pub name: String,
pub size: u64,
#[serde(default, skip_serializing_if = "is_zero")]
pub offset: u64,
pub mime_type: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub sha256: Option<String>,
#[serde(default, skip_serializing_if = "BTreeMap::is_empty")]
pub metadata: BTreeMap<String, String>,
}
impl FileTransferHeader {
pub fn new(name: impl Into<String>, size: u64, mime_type: impl Into<String>) -> Self {
Self {
version: 1,
transfer_id: uuid::Uuid::new_v4().to_string(),
name: name.into(),
size,
offset: 0,
mime_type: mime_type.into(),
sha256: None,
metadata: BTreeMap::new(),
}
}
pub fn validate(&self, limits: TransferLimits) -> Result<(), TransferError> {
if self.version != 1 {
return Err(TransferError::Protocol(
"unsupported file-transfer version".into(),
));
}
if self.transfer_id.trim().is_empty() {
return Err(TransferError::Protocol("transferId is required".into()));
}
if self.name.trim().is_empty() || self.name.contains('\0') {
return Err(TransferError::Protocol("file name is invalid".into()));
}
if self.size > limits.max_file_bytes {
return Err(TransferError::Limit(format!(
"declared file size {} exceeds {} bytes",
self.size, limits.max_file_bytes
)));
}
if self.offset > self.size {
return Err(TransferError::Protocol(
"resume offset exceeds declared file size".into(),
));
}
if let Some(digest) = &self.sha256 {
if digest.len() != 64 || !digest.bytes().all(|byte| byte.is_ascii_hexdigit()) {
return Err(TransferError::Protocol("invalid SHA-256 digest".into()));
}
}
Ok(())
}
}
#[derive(Debug, Clone, Copy)]
pub struct TransferLimits {
pub max_header_bytes: usize,
pub max_file_bytes: u64,
pub chunk_bytes: usize,
}
impl Default for TransferLimits {
fn default() -> Self {
Self {
max_header_bytes: DEFAULT_MAX_HEADER_BYTES,
max_file_bytes: DEFAULT_MAX_FILE_BYTES,
chunk_bytes: DEFAULT_CHUNK_BYTES,
}
}
}
#[derive(Debug)]
pub enum TransferError {
Io(std::io::Error),
Json(serde_json::Error),
Protocol(String),
Limit(String),
Rejected,
Integrity,
}
impl std::fmt::Display for TransferError {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Io(error) => write!(f, "file-transfer IO error: {error}"),
Self::Json(error) => write!(f, "file-transfer JSON error: {error}"),
Self::Protocol(error) => write!(f, "file-transfer protocol error: {error}"),
Self::Limit(error) => write!(f, "file-transfer limit error: {error}"),
Self::Rejected => write!(f, "file transfer rejected by receiver"),
Self::Integrity => write!(f, "file-transfer SHA-256 verification failed"),
}
}
}
impl std::error::Error for TransferError {}
impl From<std::io::Error> for TransferError {
fn from(value: std::io::Error) -> Self {
Self::Io(value)
}
}
impl From<serde_json::Error> for TransferError {
fn from(value: serde_json::Error) -> Self {
Self::Json(value)
}
}
#[derive(Debug, Clone)]
pub struct ReceivedFile {
pub header: FileTransferHeader,
pub path: PathBuf,
}
pub async fn send_path<S, F>(
stream: &mut S,
path: &Path,
mut header: FileTransferHeader,
limits: TransferLimits,
mut on_progress: F,
) -> Result<(), TransferError>
where
S: AsyncRead + AsyncWrite + Unpin,
F: FnMut(u64, u64),
{
if limits.chunk_bytes == 0 {
return Err(TransferError::Limit("chunk size must be positive".into()));
}
let mut file = tokio::fs::File::open(path).await?;
let size = file.metadata().await?.len();
header.size = size;
if header.name.trim().is_empty() {
header.name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("transfer.bin")
.to_string();
}
if header.sha256.is_some() {
header.sha256 = Some(hash_file(&mut file, limits.chunk_bytes).await?);
file.rewind().await?;
}
header.validate(limits)?;
write_header(stream, &header, limits.max_header_bytes).await?;
file.seek(std::io::SeekFrom::Start(header.offset)).await?;
let mut buffer = vec![0_u8; limits.chunk_bytes];
let mut sent = header.offset;
on_progress(sent, size);
while sent < size {
let read = file.read(&mut buffer).await?;
if read == 0 {
return Err(TransferError::Protocol(
"source file ended before its metadata size".into(),
));
}
stream.write_all(&buffer[..read]).await?;
sent += read as u64;
on_progress(sent, size);
}
stream.flush().await?;
let ack = stream.read_u8().await?;
match ack {
ACK_OK => Ok(()),
ACK_REJECTED => Err(TransferError::Rejected),
_ => Err(TransferError::Protocol(
"invalid file-transfer acknowledgement".into(),
)),
}
}
pub async fn receive_to_directory<S, F>(
stream: &mut S,
directory: &Path,
limits: TransferLimits,
mut on_progress: F,
) -> Result<ReceivedFile, TransferError>
where
S: AsyncRead + AsyncWrite + Unpin,
F: FnMut(u64, u64),
{
let result = receive_to_directory_inner(stream, directory, limits, &mut on_progress).await;
if result.is_err() {
let _ = stream.write_u8(ACK_REJECTED).await;
let _ = stream.flush().await;
}
result
}
async fn receive_to_directory_inner<S, F>(
stream: &mut S,
directory: &Path,
limits: TransferLimits,
on_progress: &mut F,
) -> Result<ReceivedFile, TransferError>
where
S: AsyncRead + AsyncWrite + Unpin,
F: FnMut(u64, u64),
{
if limits.chunk_bytes == 0 {
return Err(TransferError::Limit("chunk size must be positive".into()));
}
let header = read_header(stream, limits).await?;
tokio::fs::create_dir_all(directory).await?;
let safe_name = sanitize_file_name(&header.name);
let destination = if header.offset > 0 {
directory.join(&safe_name)
} else {
unique_destination(directory, &safe_name).await?
};
let mut hasher = header.sha256.as_ref().map(|_| Sha256::new());
if header.offset > 0 {
let metadata = tokio::fs::metadata(&destination)
.await
.map_err(|_| TransferError::Protocol("resume destination does not exist".into()))?;
if metadata.len() != header.offset {
return Err(TransferError::Protocol(format!(
"resume destination length {} does not match offset {}",
metadata.len(),
header.offset
)));
}
if let Some(hasher) = hasher.as_mut() {
let mut existing = tokio::fs::File::open(&destination).await?;
let mut prefix = vec![0_u8; limits.chunk_bytes];
loop {
let read = existing.read(&mut prefix).await?;
if read == 0 {
break;
}
hasher.update(&prefix[..read]);
}
}
}
let mut file = if header.offset > 0 {
tokio::fs::OpenOptions::new()
.append(true)
.open(&destination)
.await?
} else {
tokio::fs::OpenOptions::new()
.write(true)
.create_new(true)
.open(&destination)
.await?
};
let mut buffer = vec![0_u8; limits.chunk_bytes];
let mut received = header.offset;
on_progress(received, header.size);
while received < header.size {
let remaining = (header.size - received).min(buffer.len() as u64) as usize;
stream.read_exact(&mut buffer[..remaining]).await?;
file.write_all(&buffer[..remaining]).await?;
if let Some(hasher) = hasher.as_mut() {
hasher.update(&buffer[..remaining]);
}
received += remaining as u64;
on_progress(received, header.size);
}
file.flush().await?;
if let (Some(expected), Some(hasher)) = (&header.sha256, hasher) {
let actual = hex_digest(hasher.finalize().as_slice());
if !actual.eq_ignore_ascii_case(expected) {
drop(file);
let _ = tokio::fs::remove_file(&destination).await;
return Err(TransferError::Integrity);
}
}
stream.write_u8(ACK_OK).await?;
stream.flush().await?;
Ok(ReceivedFile {
header,
path: destination,
})
}
async fn write_header<S: AsyncWrite + Unpin>(
stream: &mut S,
header: &FileTransferHeader,
max_header_bytes: usize,
) -> Result<(), TransferError> {
let encoded = serde_json::to_vec(header)?;
if encoded.is_empty() || encoded.len() > max_header_bytes || encoded.len() > u32::MAX as usize {
return Err(TransferError::Limit(format!(
"header exceeds {max_header_bytes} bytes"
)));
}
stream.write_u32(encoded.len() as u32).await?;
stream.write_all(&encoded).await?;
Ok(())
}
async fn read_header<S: AsyncRead + Unpin>(
stream: &mut S,
limits: TransferLimits,
) -> Result<FileTransferHeader, TransferError> {
let length = stream.read_u32().await? as usize;
if length == 0 || length > limits.max_header_bytes {
return Err(TransferError::Limit(format!(
"invalid header length {length}"
)));
}
let mut encoded = vec![0_u8; length];
stream.read_exact(&mut encoded).await?;
let header: FileTransferHeader = serde_json::from_slice(&encoded)?;
header.validate(limits)?;
Ok(header)
}
async fn hash_file(
file: &mut tokio::fs::File,
chunk_bytes: usize,
) -> Result<String, TransferError> {
let mut hasher = Sha256::new();
let mut buffer = vec![0_u8; chunk_bytes];
loop {
let read = file.read(&mut buffer).await?;
if read == 0 {
break;
}
hasher.update(&buffer[..read]);
}
Ok(hex_digest(hasher.finalize().as_slice()))
}
fn hex_digest(bytes: &[u8]) -> String {
let mut output = String::with_capacity(bytes.len() * 2);
for byte in bytes {
use std::fmt::Write as _;
let _ = write!(output, "{byte:02x}");
}
output
}
fn is_zero(value: &u64) -> bool {
*value == 0
}
pub fn sanitize_file_name(value: &str) -> String {
let candidate = value.replace('\\', "/");
let name = candidate.rsplit('/').next().unwrap_or("").trim();
let sanitized: String = name
.chars()
.map(|character| {
if character.is_control() || character == '\0' {
'_'
} else {
character
}
})
.take(240)
.collect();
match sanitized.as_str() {
"" | "." | ".." => "transfer.bin".to_string(),
_ => sanitized,
}
}
async fn unique_destination(directory: &Path, file_name: &str) -> Result<PathBuf, TransferError> {
let original = Path::new(file_name);
let stem = original
.file_stem()
.and_then(|value| value.to_str())
.unwrap_or("transfer");
let extension = original.extension().and_then(|value| value.to_str());
for suffix in 0..10_000_u32 {
let name = if suffix == 0 {
file_name.to_string()
} else if let Some(extension) = extension {
format!("{stem} ({suffix}).{extension}")
} else {
format!("{stem} ({suffix})")
};
let candidate = directory.join(name);
if !tokio::fs::try_exists(&candidate).await? {
return Ok(candidate);
}
}
Err(TransferError::Limit(
"unable to allocate a unique destination name".into(),
))
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn round_trip_streams_and_verifies_a_file() {
let source_dir = tempfile::tempdir().unwrap();
let destination_dir = tempfile::tempdir().unwrap();
let source = source_dir.path().join("plane.webp");
tokio::fs::write(&source, [1_u8, 2, 3, 4, 5]).await.unwrap();
let (mut sender, mut receiver) = tokio::io::duplex(256 * 1024);
let mut header = FileTransferHeader::new("plane.webp", 0, "image/webp");
header.transfer_id = "transfer-1".to_string();
header.sha256 = Some(String::new());
let limits = TransferLimits::default();
let send = send_path(&mut sender, &source, header, limits, |_, _| {});
let receive =
receive_to_directory(&mut receiver, destination_dir.path(), limits, |_, _| {});
let (sent, received) = tokio::join!(send, receive);
sent.unwrap();
let received = received.unwrap();
assert_eq!(received.header.transfer_id, "transfer-1");
assert_eq!(
tokio::fs::read(received.path).await.unwrap(),
[1, 2, 3, 4, 5]
);
}
#[test]
fn sanitizes_cross_platform_path_traversal() {
assert_eq!(sanitize_file_name("../../secret.txt"), "secret.txt");
assert_eq!(
sanitize_file_name(r"C:\\Users\\name\\secret.txt"),
"secret.txt"
);
assert_eq!(sanitize_file_name(".."), "transfer.bin");
}
#[tokio::test]
async fn rejects_oversized_headers_before_creating_a_file() {
let destination = tempfile::tempdir().unwrap();
let (mut sender, mut receiver) = tokio::io::duplex(1024);
sender.write_u32(1024).await.unwrap();
sender.flush().await.unwrap();
let limits = TransferLimits {
max_header_bytes: 16,
..TransferLimits::default()
};
let error = receive_to_directory(&mut receiver, destination.path(), limits, |_, _| {})
.await
.unwrap_err();
assert!(matches!(error, TransferError::Limit(_)));
assert_eq!(std::fs::read_dir(destination.path()).unwrap().count(), 0);
}
#[tokio::test]
async fn resumes_an_existing_verified_prefix() {
let source_dir = tempfile::tempdir().unwrap();
let destination_dir = tempfile::tempdir().unwrap();
let source = source_dir.path().join("resume.bin");
let destination = destination_dir.path().join("resume.bin");
tokio::fs::write(&source, [1_u8, 2, 3, 4, 5]).await.unwrap();
tokio::fs::write(&destination, [1_u8, 2]).await.unwrap();
let (mut sender, mut receiver) = tokio::io::duplex(256 * 1024);
let mut header = FileTransferHeader::new("resume.bin", 0, "application/octet-stream");
header.transfer_id = "resume-1".to_string();
header.offset = 2;
header.sha256 = Some(String::new());
let limits = TransferLimits::default();
let send = send_path(&mut sender, &source, header, limits, |_, _| {});
let receive =
receive_to_directory(&mut receiver, destination_dir.path(), limits, |_, _| {});
let (sent, received) = tokio::join!(send, receive);
sent.unwrap();
assert_eq!(received.unwrap().header.offset, 2);
assert_eq!(tokio::fs::read(destination).await.unwrap(), [1, 2, 3, 4, 5]);
}
#[test]
fn routed_frames_round_trip_offsets_chunks_and_completion() {
let limits = TransferLimits {
chunk_bytes: DEFAULT_ROUTED_CHUNK_BYTES,
..TransferLimits::default()
};
let mut header = FileTransferHeader::new("route.bin", 99, "application/octet-stream");
header.transfer_id = "route-transfer-1".to_string();
let frames = vec![
RoutedTransferFrame::Offer(header),
RoutedTransferFrame::Resume {
transfer_id: "route-transfer-1".to_string(),
offset: 24,
},
RoutedTransferFrame::Chunk {
transfer_id: "route-transfer-1".to_string(),
offset: 24,
ack_requested: true,
bytes: vec![1, 2, 3],
},
RoutedTransferFrame::Ack {
transfer_id: "route-transfer-1".to_string(),
offset: 27,
},
RoutedTransferFrame::Complete {
transfer_id: "route-transfer-1".to_string(),
size: 99,
},
RoutedTransferFrame::Completed {
transfer_id: "route-transfer-1".to_string(),
size: 99,
},
RoutedTransferFrame::Reject {
transfer_id: "route-transfer-1".to_string(),
reason: "ACL changed".to_string(),
},
];
for frame in frames {
let encoded = frame.encode(limits).unwrap();
assert_eq!(
RoutedTransferFrame::decode(&encoded, limits).unwrap(),
frame
);
}
}
#[test]
fn routed_chunk_rejects_oversized_payloads() {
let limits = TransferLimits {
chunk_bytes: 4,
..TransferLimits::default()
};
let frame = RoutedTransferFrame::Chunk {
transfer_id: "route-transfer-1".to_string(),
offset: 0,
ack_requested: true,
bytes: vec![0; 5],
};
assert!(matches!(frame.encode(limits), Err(TransferError::Limit(_))));
}
}