use crate::crypto::{Sha256PrefixState, TokenOpenStream};
use crate::routing::links::resources::assemble_incoming::{
verify_absorbed_and_prove, verify_and_prove, OpenTransferError, VerifyResourceError,
};
use crate::routing::links::resources::{
ResourceCompression, ResourceHash, ResourceProof, SaltNonce, RESOURCE_NONCE_LEN,
};
use crate::routing::links::LinkKey;
pub struct StreamedOpen {
token: TokenOpenStream,
stream_digest: Option<Sha256PrefixState>,
plaintext_seen_byte_len: usize,
}
impl core::fmt::Debug for StreamedOpen {
fn fmt(&self, f: &mut core::fmt::Formatter<'_>) -> core::fmt::Result {
f.debug_struct("StreamedOpen").finish_non_exhaustive()
}
}
impl StreamedOpen {
pub fn begin(key: &LinkKey, sealed: &[u8], compression: ResourceCompression) -> Option<Self> {
let iv = sealed.get(..16)?.try_into().ok()?;
let token = key.open_stream(iv, sealed.len()).ok()?;
let stream_digest = match compression {
ResourceCompression::Uncompressed => Some(Sha256PrefixState::absorb(&[])),
ResourceCompression::Bz2 => None,
};
Some(Self {
token,
stream_digest,
plaintext_seen_byte_len: 0,
})
}
pub fn advance(&mut self, transfer: &mut [u8], contiguous_byte_len: usize) {
let span = self.token.pending_span(contiguous_byte_len);
self.chew_span(&mut transfer[span]);
}
pub fn pending_span(&self, contiguous_byte_len: usize) -> core::ops::Range<usize> {
self.token.pending_span(contiguous_byte_len)
}
pub fn caught_up(&self) -> bool {
self.token.fully_absorbed()
}
pub fn chew_span(&mut self, span: &mut [u8]) {
self.token.absorb_span(span);
let nonce_still_owed = RESOURCE_NONCE_LEN.saturating_sub(self.plaintext_seen_byte_len);
self.plaintext_seen_byte_len += span.len();
let Some(digest) = &mut self.stream_digest else {
return;
};
digest.update(&span[nonce_still_owed.min(span.len())..]);
}
pub fn conclude(mut self, transfer: &mut [u8]) -> Result<OpenedStream<'_>, OpenTransferError> {
self.advance(transfer, transfer.len());
let plaintext = self
.token
.finalize(transfer)
.map_err(OpenTransferError::Open)?;
if plaintext.len() < RESOURCE_NONCE_LEN {
return Err(OpenTransferError::StreamTooShort);
}
let stream = &plaintext[RESOURCE_NONCE_LEN..];
let absorbed = self.stream_digest.map(|mut digest| {
let digested = self
.plaintext_seen_byte_len
.saturating_sub(RESOURCE_NONCE_LEN)
.min(stream.len());
digest.update(&stream[digested..]);
digest
});
Ok(OpenedStream { stream, absorbed })
}
}
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Default)]
pub enum OpenProgress {
#[default]
NotBegun,
Parked(StreamedOpen),
Chewing {
dispatched: core::ops::Range<usize>,
},
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum ResourceOpenLane {
#[default]
Inline,
PoolWhenContended,
}
pub struct OpenedStream<'t> {
pub stream: &'t [u8],
absorbed: Option<Sha256PrefixState>,
}
impl<'t> OpenedStream<'t> {
pub fn rehashing(stream: &'t [u8]) -> Self {
Self {
stream,
absorbed: None,
}
}
pub fn verify_and_prove(
&self,
salt_nonce: &SaltNonce,
advertised: &ResourceHash,
) -> Result<ResourceProof, VerifyResourceError> {
match &self.absorbed {
Some(midstate) => verify_absorbed_and_prove(midstate, salt_nonce, advertised),
None => verify_and_prove(self.stream, salt_nonce, advertised),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::crypto::{sha256, x25519_diffie_hellman, X25519PublicKey, X25519SecretKey};
use crate::routing::links::resources::build_outgoing::{build_outgoing_resource, BuildRegions};
use crate::routing::links::resources::{resource_sdu, ResourceBody, ResourceMetadata};
use crate::routing::links::LinkId;
use crate::wire::BROADCAST_MTU;
fn link_key() -> LinkKey {
let shared = x25519_diffie_hellman(
&X25519SecretKey::new([0x33; 32]),
&X25519PublicKey([0x55; 32]),
);
LinkKey::derive(&LinkId::new([0x07; 16]), &shared)
}
fn nonces() -> impl FnMut() -> [u8; RESOURCE_NONCE_LEN] {
let mut drawn = 0;
move || {
drawn += 1;
if drawn == 1 {
[0x51, 0x52, 0x53, 0x54]
} else {
[0x61, 0x62, 0x63, 0x64]
}
}
}
fn payload() -> std::vec::Vec<u8> {
let mut seed = sha256(b"streamed-open");
let mut data = std::vec::Vec::new();
for _ in 0..47 {
data.extend_from_slice(&seed);
seed = sha256(&seed);
}
data.truncate(1_400);
data
}
fn built_transfer(
body: &ResourceBody<'_>,
) -> (
std::vec::Vec<u8>,
crate::routing::links::resources::build_outgoing::BuiltResource,
usize,
) {
let mut transfer = std::vec![0u8; 4_096];
let mut hashmap = std::vec![0u8; 64];
let sdu = resource_sdu(BROADCAST_MTU);
let built = build_outgoing_resource(
body,
&link_key(),
&[0xA1; 16],
nonces(),
sdu,
BuildRegions {
transfer: &mut transfer,
hashmap: &mut hashmap,
},
)
.unwrap();
transfer.truncate(built.sealed_transfer_bytes);
(transfer, built, sdu)
}
#[test]
fn an_uncompressed_transfer_streamed_part_by_part_verifies_off_the_midstate() {
let data = payload();
let (mut transfer, built, sdu) = built_transfer(&ResourceBody {
data: &data,
compressed_candidate: None,
metadata: ResourceMetadata::None,
});
let mut open =
StreamedOpen::begin(&link_key(), &transfer, ResourceCompression::Uncompressed).unwrap();
for part_index in 0..built.part_count {
let contiguous = ((part_index + 1) * sdu).min(transfer.len());
open.advance(&mut transfer, contiguous);
}
let opened = open.conclude(&mut transfer).unwrap();
assert_eq!(opened.stream, &data[..]);
assert_eq!(
opened
.verify_and_prove(&built.salt_nonce, &built.hash)
.unwrap(),
built.expected_proof,
);
}
#[test]
fn a_stalled_frontier_that_jumps_at_the_end_still_verifies() {
let data = payload();
let (mut transfer, built, _) = built_transfer(&ResourceBody {
data: &data,
compressed_candidate: None,
metadata: ResourceMetadata::None,
});
let mut open =
StreamedOpen::begin(&link_key(), &transfer, ResourceCompression::Uncompressed).unwrap();
open.advance(&mut transfer, 100);
open.advance(&mut transfer, 100);
let opened = open.conclude(&mut transfer).unwrap();
assert_eq!(opened.stream, &data[..]);
assert_eq!(
opened
.verify_and_prove(&built.salt_nonce, &built.hash)
.unwrap(),
built.expected_proof,
);
}
#[test]
fn a_compressed_transfer_streams_the_decrypt_and_yields_the_bz2_stream() {
let data = payload();
let candidate = b"pretend bz2, just visibly shorter".to_vec();
let (mut transfer, built, sdu) = built_transfer(&ResourceBody {
data: &data,
compressed_candidate: Some(&candidate),
metadata: ResourceMetadata::None,
});
let mut open =
StreamedOpen::begin(&link_key(), &transfer, ResourceCompression::Bz2).unwrap();
let contiguous = sdu.min(transfer.len());
open.advance(&mut transfer, contiguous);
let opened = open.conclude(&mut transfer).unwrap();
assert_eq!(opened.stream, &candidate[..]);
assert_eq!(built.compression, ResourceCompression::Bz2);
}
#[test]
fn a_tampered_transfer_refuses_to_conclude() {
let data = payload();
let (mut transfer, _, _) = built_transfer(&ResourceBody {
data: &data,
compressed_candidate: None,
metadata: ResourceMetadata::None,
});
*transfer.last_mut().unwrap() ^= 1;
let mut open =
StreamedOpen::begin(&link_key(), &transfer, ResourceCompression::Uncompressed).unwrap();
open.advance(&mut transfer, usize::MAX);
assert!(matches!(
open.conclude(&mut transfer),
Err(OpenTransferError::Open(_)),
));
}
#[test]
fn a_buffer_too_short_for_any_token_never_begins() {
assert!(
StreamedOpen::begin(&link_key(), &[0u8; 12], ResourceCompression::Uncompressed,)
.is_none()
);
assert!(
StreamedOpen::begin(&link_key(), &[0u8; 63], ResourceCompression::Uncompressed,)
.is_none()
);
}
}