crime 0.6.1

Concurrent real-time interface for multimedia engines
Documentation
#[cfg(feature = "webm")]
use crate::webm;
use crate::{OpusApplication, OpusBitrate};
use futures::{Stream, StreamExt};
use opus::Channels as OpusChannels;
use std::pin::Pin;

pub struct OpusPacket {
    pub data: Vec<u8>,
    pub frame_size_samples: usize,
}

pub async fn encode_opus_stream<'a>(
    samples: impl Stream<Item = f32> + Send + 'a,
    sample_rate: u32,
    application: OpusApplication,
    bitrate: OpusBitrate,
) -> impl Stream<Item = Result<OpusPacket, opus::Error>> + Send + 'a {
    // Opus supports 2.5, 5, 10, 20, 40, 60 ms.
    // We use 20ms as standard frame size.
    let frame_size = (sample_rate as usize * 20) / 1000;
    let samples = Box::pin(samples);
    let mut sample_chunks = samples.chunks(frame_size);

    async_stream::stream! {
        // Initialize encoder inside the stream so errors can be yielded
        let mut encoder = match opus::Encoder::new(sample_rate, OpusChannels::Mono, application) {
            Ok(enc) => enc,
            Err(e) => {
                yield Err(e);
                return;
            }
        };

        if let Err(e) = encoder.set_bitrate(bitrate) {
            yield Err(e);
            return;
        }

        while let Some(chunk) = sample_chunks.next().await {
            let mut output = [0u8; 4000];

            // Handle partial frames by padding
            let encode_result = if chunk.len() == frame_size {
                encoder.encode_float(&chunk, &mut output)
            } else {
                let mut padded = chunk;
                padded.resize(frame_size, 0.0);
                encoder.encode_float(&padded, &mut output)
            };

            match encode_result {
                Ok(len) => {
                    yield Ok(OpusPacket {
                        data: output[..len].to_vec(),
                        frame_size_samples: frame_size,
                    });
                }
                Err(e) => {
                    yield Err(e);
                    return;
                }
            }
        }
    }
}

#[cfg(any(feature = "ogg", feature = "webm"))]
const OPUS_HEAD_MAGIC: &[u8] = b"OpusHead";
#[cfg(feature = "ogg")]
const OPUS_TAGS_MAGIC: &[u8] = b"OpusTags";

pub struct OpusHeader {
    pub version: u8,
    pub channels: u8,
    pub pre_skip: u16,
    pub input_sample_rate: u32,
    pub output_gain: i16,
    pub channel_mapping_family: u8,
    // Optional: channel mapping table
}

impl OpusHeader {
    pub fn to_bytes(&self) -> Vec<u8> {
        let mut bytes = Vec::new();
        bytes.extend_from_slice(OPUS_HEAD_MAGIC);
        bytes.push(self.version);
        bytes.push(self.channels);
        bytes.extend_from_slice(&self.pre_skip.to_le_bytes());
        bytes.extend_from_slice(&self.input_sample_rate.to_le_bytes());
        bytes.extend_from_slice(&self.output_gain.to_le_bytes());
        bytes.push(self.channel_mapping_family);
        // Mapping family 0 implies mono or stereo (L, R) - no table needed
        bytes
    }
}

#[cfg(feature = "ogg")]
pub fn make_opus_comment_header() -> Vec<u8> {
    let mut bytes = Vec::new();
    bytes.extend_from_slice(OPUS_TAGS_MAGIC);

    // Vendor String Length (u32 le)
    let vendor = "rust-crime-crate";
    bytes.extend_from_slice(&(vendor.len() as u32).to_le_bytes());
    bytes.extend_from_slice(vendor.as_bytes());

    // User Comment List Length (u32 le) - 0 for now
    bytes.extend_from_slice(&0u32.to_le_bytes());

    bytes
}

#[cfg(feature = "ogg")]
pub async fn encode_opus_as_ogg<'a>(
    samples: impl Stream<Item = f32> + Send + 'a,
    sample_rate: u32,
    application: OpusApplication,
    bitrate: OpusBitrate,
) -> Pin<Box<dyn Stream<Item = Vec<u8>> + Send + 'a>> {
    let opus_packets = encode_opus_stream(samples, sample_rate, application, bitrate).await;
    let mut opus_packets = Box::pin(opus_packets);

    Box::pin(async_stream::stream! {
        const OGG_CRC_ALGO: crc::Algorithm<u32> = crc::Algorithm {
            width: 32,
            poly: 0x04c11db7,
            init: 0,
            refin: false,
            refout: false,
            xorout: 0,
            check: 0,
            residue: 0,
        };
        let ogg_crc = crc::Crc::<u32>::new(&OGG_CRC_ALGO);

        // ID Header Page
        let id_header = OpusHeader {
            version: 1,
            channels: 1,
            pre_skip: 0,
            input_sample_rate: sample_rate,
            output_gain: 0,
            channel_mapping_family: 0,
        };
        let id_packet = id_header.to_bytes();

        let mut id_page = Vec::new();
        id_page.extend_from_slice(b"OggS");
        id_page.push(0); // version
        id_page.push(0x02); // type: BOS
        id_page.extend_from_slice(&0u64.to_le_bytes()); // granule pos
        id_page.extend_from_slice(&1u32.to_le_bytes()); // serial
        id_page.extend_from_slice(&0u32.to_le_bytes()); // sequence
        id_page.extend_from_slice(&0u32.to_le_bytes()); // checksum placeholder
        id_page.push(1); // segments
        id_page.push(id_packet.len() as u8);
        id_page.extend_from_slice(&id_packet);

        let crc = ogg_crc.checksum(&id_page);
        id_page[22..26].copy_from_slice(&crc.to_le_bytes());

        // Comment Header Page
        let comment_packet = make_opus_comment_header();
        let mut comment_page = Vec::new();
        comment_page.extend_from_slice(b"OggS");
        comment_page.push(0); // version
        comment_page.push(0); // type: normal
        comment_page.extend_from_slice(&0u64.to_le_bytes()); // granule pos
        comment_page.extend_from_slice(&1u32.to_le_bytes()); // serial
        comment_page.extend_from_slice(&1u32.to_le_bytes()); // sequence
        comment_page.extend_from_slice(&0u32.to_le_bytes()); // checksum placeholder
        comment_page.push(1); // segments
        comment_page.push(comment_packet.len() as u8);
        comment_page.extend_from_slice(&comment_packet);

        let crc = ogg_crc.checksum(&comment_page);
        comment_page[22..26].copy_from_slice(&crc.to_le_bytes());

        // Build preamble (ID + Comment pages) to prepend to first audio page.
        let mut preamble = Vec::with_capacity(id_page.len() + comment_page.len());
        preamble.extend_from_slice(&id_page);
        preamble.extend_from_slice(&comment_page);

        // Audio Pages
        let mut granule_position: u64 = 0;
        let mut page_sequence: u32 = 2;
        let mut first = true;

        while let Some(packet_result) = opus_packets.next().await {
            let packet = match packet_result {
                Ok(p) => p,
                Err(e) => {
                    eprintln!("Opus encoding error: {:?}", e);
                    return;
                }
            };
            let opus_data = packet.data;
            granule_position += packet.frame_size_samples as u64;

            let mut page = Vec::new();
            page.extend_from_slice(b"OggS");
            page.push(0); // version
            page.push(0); // type: normal

            page.extend_from_slice(&granule_position.to_le_bytes());
            page.extend_from_slice(&1u32.to_le_bytes()); // serial
            page.extend_from_slice(&page_sequence.to_le_bytes());
            page.extend_from_slice(&0u32.to_le_bytes()); // checksum placeholder
            page.push(1); // segments placeholder

            let mut len_remaining = opus_data.len();
            while len_remaining >= 255 {
                page.push(255);
                len_remaining -= 255;
            }
            page.push(len_remaining as u8);

            // 27 = fixed header size (up to and including the segment count byte)
            let num_segments = page.len() - 27;
            page[26] = num_segments as u8;

            page.extend_from_slice(&opus_data);

            let crc = ogg_crc.checksum(&page);
            page[22..26].copy_from_slice(&crc.to_le_bytes());

            if first {
                first = false;
                let mut buf = Vec::with_capacity(preamble.len() + page.len());
                buf.extend_from_slice(&preamble);
                buf.extend_from_slice(&page);
                yield buf;
            } else {
                yield page;
            }

            page_sequence += 1;
        }

        // If no audio packets, yield preamble alone.
        if first {
            yield preamble;
        }
    })
}

#[cfg(feature = "webm")]
pub fn make_webm_header() -> Vec<u8> {
    let mut header = Vec::new();
    header.extend_from_slice(&webm::make_uint_element(0x4286, 1)); // EBMLVersion
    header.extend_from_slice(&webm::make_uint_element(0x42F7, 1)); // EBMLReadVersion
    header.extend_from_slice(&webm::make_uint_element(0x42F2, 4)); // EBMLMaxIDLength
    header.extend_from_slice(&webm::make_uint_element(0x42F3, 8)); // EBMLMaxSizeLength
    header.extend_from_slice(&webm::make_string_element(0x4282, "webm")); // DocType
    header.extend_from_slice(&webm::make_uint_element(0x4287, 4)); // DocTypeVersion
    header.extend_from_slice(&webm::make_uint_element(0x4285, 2)); // DocTypeReadVersion
    webm::make_element(webm::EBML_ID, &header)
}

#[cfg(feature = "webm")]
pub fn make_segment_header() -> Vec<u8> {
    let mut segment = Vec::new();
    segment.extend_from_slice(&webm::encode_id(webm::SEGMENT_ID));
    // Unknown size (all 1s, 8 bytes width -> 0x01FFFFFFFFFFFFFF)
    let unknown_size = [0x01, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF, 0xFF];
    segment.extend_from_slice(&unknown_size);
    segment
}

#[cfg(feature = "webm")]
pub fn make_info_element() -> Vec<u8> {
    let mut info = Vec::new();
    info.extend_from_slice(&webm::make_uint_element(webm::TIMECODE_SCALE_ID, 1_000_000)); // 1ms
    let app_string = concat!(env!("CARGO_PKG_NAME"), " ", env!("CARGO_PKG_VERSION"));
    info.extend_from_slice(&webm::make_string_element(webm::MUXING_APP_ID, app_string));
    info.extend_from_slice(&webm::make_string_element(webm::WRITING_APP_ID, app_string));
    webm::make_element(webm::INFO_ID, &info)
}

#[cfg(feature = "webm")]
pub fn make_tracks_element(sample_rate: u32) -> Vec<u8> {
    let mut tracks = Vec::new();
    let mut track_entry = Vec::new();
    track_entry.extend_from_slice(&webm::make_uint_element(webm::TRACK_NUMBER_ID, 1));
    track_entry.extend_from_slice(&webm::make_uint_element(webm::TRACK_UID_ID, 1));
    track_entry.extend_from_slice(&webm::make_uint_element(webm::TRACK_TYPE_ID, 2)); // Audio
    track_entry.extend_from_slice(&webm::make_string_element(webm::CODEC_ID_ID, "A_OPUS"));

    // CodecDelay (Opus Pre-Skip): 312 samples at 48k = 6.5ms
    let pre_skip = ((sample_rate as u64 * 312) / 48000) as u16;
    let codec_delay = 6_500_000u64; // 6.5ms in ns
    let seek_pre_roll = 80_000_000u64; // 80ms in ns

    track_entry.extend_from_slice(&webm::make_uint_element(webm::CODEC_DELAY_ID, codec_delay));
    track_entry.extend_from_slice(&webm::make_uint_element(
        webm::SEEK_PRE_ROLL_ID,
        seek_pre_roll,
    ));

    // CodecPrivate: OpusHead
    let opus_head = OpusHeader {
        version: 1,
        channels: 1,
        pre_skip,
        input_sample_rate: sample_rate,
        output_gain: 0,
        channel_mapping_family: 0,
    };
    track_entry.extend_from_slice(&webm::make_element(
        webm::CODEC_PRIVATE_ID,
        &opus_head.to_bytes(),
    ));

    // Audio sub-element
    let mut audio = Vec::new();
    audio.extend_from_slice(&webm::make_float_element(
        webm::SAMPLING_FREQUENCY_ID,
        sample_rate as f32,
    ));
    audio.extend_from_slice(&webm::make_uint_element(webm::CHANNELS_ID, 1));
    track_entry.extend_from_slice(&webm::make_element(webm::AUDIO_ID, &audio));

    tracks.extend_from_slice(&webm::make_element(webm::TRACK_ENTRY_ID, &track_entry));

    webm::make_element(webm::TRACKS_ID, &tracks)
}

#[cfg(feature = "webm")]
pub async fn encode_opus_as_webm<'a>(
    samples: impl Stream<Item = f32> + Send + 'a,
    sample_rate: u32,
    application: OpusApplication,
    bitrate: OpusBitrate,
) -> Pin<Box<dyn Stream<Item = Vec<u8>> + Send + 'a>> {
    let opus_packets = encode_opus_stream(samples, sample_rate, application, bitrate).await;
    let mut opus_packets = Box::pin(opus_packets);

    Box::pin(async_stream::stream! {
        // Build the full WebM header (EBML + Segment + Info + Tracks) and
        // prepend it to the first audio cluster so the first yielded chunk
        // always contains audio data.
        let mut preamble = Vec::new();
        preamble.extend_from_slice(&make_webm_header());
        preamble.extend_from_slice(&make_segment_header());
        preamble.extend_from_slice(&make_info_element());
        preamble.extend_from_slice(&make_tracks_element(sample_rate));

        let mut cluster_timecode = 0u64;
        let mut first = true;

        while let Some(packet_result) = opus_packets.next().await {
            let packet = match packet_result {
                Ok(p) => p,
                Err(e) => {
                    eprintln!("Opus encoding error: {:?}", e);
                    return;
                }
            };

            let mut cluster_data = Vec::new();
            cluster_data.extend_from_slice(&webm::make_uint_element(webm::TIMECODE_ID, cluster_timecode));
            cluster_data.extend_from_slice(&webm::make_simple_block(1, 0, &packet.data));

            let cluster = webm::make_element(webm::CLUSTER_ID, &cluster_data);

            if first {
                first = false;
                let mut buf = Vec::with_capacity(preamble.len() + cluster.len());
                buf.extend_from_slice(&preamble);
                buf.extend_from_slice(&cluster);
                yield buf;
            } else {
                yield cluster;
            }

            cluster_timecode += 20; // 20ms per frame
        }

        // If no audio packets were produced, yield the preamble alone.
        if first {
            yield preamble;
        }
    })
}