use crate::cache;
use crate::frame::{self, Frame, FrameBuf};
use crate::{Timescale, stats, track};
use std::collections::VecDeque;
use std::mem::MaybeUninit;
use std::sync::Arc;
use std::task::{Poll, ready};
use crate::{Error, IntoBytes, Result, Timestamp};
const MAX_GROUP_CACHE: u64 = 32 * 1024 * 1024;
#[derive(Clone, Copy, Debug, Hash, Eq, PartialEq, Ord, PartialOrd)]
pub struct Info {
pub sequence: u64,
}
impl Info {
#[cfg(test)]
pub(crate) fn produce(self) -> Producer {
Producer::new(self, track::Info::default())
}
}
impl From<usize> for Info {
fn from(sequence: usize) -> Self {
Self {
sequence: sequence as u64,
}
}
}
impl From<u64> for Info {
fn from(sequence: u64) -> Self {
Self { sequence }
}
}
impl From<u32> for Info {
fn from(sequence: u32) -> Self {
Self {
sequence: sequence as u64,
}
}
}
impl From<u16> for Info {
fn from(sequence: u16) -> Self {
Self {
sequence: sequence as u64,
}
}
}
pub(crate) struct Partial {
timestamp: Timestamp,
buf: FrameBuf,
}
#[derive(Default)]
pub(crate) struct GroupState {
pub(crate) frames: VecDeque<Frame>,
pub(crate) partial: Option<Partial>,
pub(crate) offset: usize,
pub(crate) cache: u64,
charge: cache::Charge,
pub(crate) fin: bool,
pub(crate) abort: Option<Error>,
}
impl GroupState {
fn poll_frame_source(&self, index: usize) -> Poll<Result<Option<(frame::Info, frame::Source)>>> {
if index < self.offset {
return Poll::Ready(Err(Error::Lagged));
}
let local = index - self.offset;
if let Some(f) = self.frames.get(local) {
self.charge.touch();
let info = frame::Info {
size: f.payload.len() as u64,
timestamp: f.timestamp,
};
return Poll::Ready(Ok(Some((info, frame::Source::Complete(f.payload.clone())))));
}
if local == self.frames.len()
&& let Some(p) = &self.partial
{
self.charge.touch();
let info = frame::Info {
size: p.buf.capacity() as u64,
timestamp: p.timestamp,
};
return Poll::Ready(Ok(Some((info, frame::Source::Partial(p.buf.clone())))));
}
if let Some(err) = &self.abort {
return Poll::Ready(Err(err.clone()));
}
if self.fin {
return Poll::Ready(Ok(None));
}
Poll::Pending
}
fn poll_finished(&self) -> Poll<Result<u64>> {
if let Some(err) = &self.abort {
Poll::Ready(Err(err.clone()))
} else if self.fin {
Poll::Ready(Ok((self.offset + self.frames.len()) as u64))
} else {
Poll::Pending
}
}
fn evict(&mut self) {
while self.cache > MAX_GROUP_CACHE {
let Some(frame) = self.frames.pop_front() else {
break;
};
let size = frame.payload.len() as u64;
self.cache -= size;
self.charge.sub(size);
self.offset += 1;
}
}
fn release(&mut self) {
self.frames.clear();
self.partial = None;
self.cache = 0;
self.charge.clear();
}
}
fn modify(state: &kio::Producer<GroupState>) -> Result<kio::Mut<'_, GroupState>> {
state.write().map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
}
fn evict(weak: &kio::ProducerWeak<GroupState>) {
let Some(producer) = weak.produce() else { return };
let Ok(mut state) = producer.write() else { return };
if state.abort.is_some() {
return;
}
state.abort = Some(Error::Evicted);
state.release();
state.close();
}
pub struct Producer {
state: kio::Producer<GroupState>,
info: Info,
track: track::Info,
stats: stats::Meter,
}
impl std::ops::Deref for Producer {
type Target = Info;
fn deref(&self) -> &Self::Target {
&self.info
}
}
impl Producer {
pub(crate) fn new(info: Info, track: track::Info) -> Self {
let state = kio::Producer::<GroupState>::default();
let weak = state.weak();
let charge = track.broadcast.origin.pool.register(Box::new(move || evict(&weak)));
state.write().ok().expect("a new group is open").charge = charge;
Self {
info,
state,
track,
stats: stats::Meter::default(),
}
}
pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
meter.group();
self.stats = meter;
self
}
pub(crate) fn info(&self) -> Info {
self.info
}
pub fn timescale(&self) -> Timescale {
self.track.timescale
}
pub fn write_frame<B: IntoBytes>(&mut self, timestamp: Timestamp, data: B) -> Result<()> {
let timestamp = timestamp
.convert(self.track.timescale)
.map_err(|_| Error::TimestampMismatch)?;
let payload = data.into_bytes();
if payload.len() as u64 > MAX_GROUP_CACHE {
return Err(Error::FrameTooLarge);
}
let mut state = modify(&self.state)?;
if state.fin {
return Err(Error::Closed);
}
debug_assert!(state.partial.is_none(), "a frame is already open");
let size = payload.len() as u64;
state.cache += size;
state.charge.add(size);
state.frames.push_back(Frame { timestamp, payload });
state.evict();
drop(state);
self.track.broadcast.origin.pool.evict();
self.stats.frames(1);
self.stats.bytes(size);
Ok(())
}
pub fn create_frame(&mut self, frame: frame::Info) -> Result<frame::Producer<'_>> {
let timestamp = frame
.timestamp
.convert(self.track.timescale)
.map_err(|_| Error::TimestampMismatch)?;
if frame.size > MAX_GROUP_CACHE {
return Err(Error::FrameTooLarge);
}
let buf = FrameBuf::new(frame.size as usize);
let mut state = modify(&self.state)?;
if state.fin {
return Err(Error::Closed);
}
debug_assert!(state.partial.is_none(), "a frame is already open");
state.cache += frame.size;
state.charge.add(frame.size);
state.partial = Some(Partial {
timestamp,
buf: buf.clone(),
});
state.evict();
drop(state);
self.track.broadcast.origin.pool.evict();
self.stats.frames(1);
let meter = self.stats.clone();
let info = frame::Info {
size: frame.size,
timestamp,
};
Ok(frame::Producer::new(self, buf, info).with_meter(meter))
}
pub(crate) fn frame_notify(&self) {
let _ = self.state.write();
}
pub(crate) fn frame_commit(&mut self, frame: Frame) -> Result<()> {
let mut state = modify(&self.state)?;
state.partial = None;
state.frames.push_back(frame);
Ok(())
}
pub(crate) fn frame_abort(&mut self, err: Error) {
let _ = self.clone().abort(err);
}
pub fn frame_count(&self) -> usize {
let state = self.state.read();
state.offset + state.frames.len() + state.partial.is_some() as usize
}
pub fn finish(&mut self) -> Result<()> {
let mut state = modify(&self.state)?;
state.fin = true;
Ok(())
}
pub fn abort(self, err: Error) -> Result<()> {
let mut guard = modify(&self.state)?;
guard.abort = Some(err);
guard.release();
guard.close();
Ok(())
}
pub(crate) fn is_aborted(&self) -> bool {
self.state.read().abort.is_some()
}
pub(crate) fn cache_entry(&self) -> Option<Arc<cache::Entry>> {
self.state.read().charge.entry()
}
pub fn consume(&self) -> Consumer {
Consumer {
info: self.info,
state: self.state.consume(),
track: self.track.clone(),
index: 0,
prefetch: Prefetch::default(),
stats: stats::Meter::default(),
}
}
pub async fn closed(&self) -> Error {
kio::wait(|waiter| self.poll_closed(waiter)).await
}
pub fn poll_closed(&self, waiter: &kio::Waiter) -> Poll<Error> {
self.state.poll_closed(waiter).map(|()| self.abort_reason())
}
pub async fn unused(&self) -> Result<()> {
self.state.unused().await.map_err(|_| self.abort_reason())
}
fn abort_reason(&self) -> Error {
self.state.read().abort.clone().unwrap_or(Error::Dropped)
}
}
impl Clone for Producer {
fn clone(&self) -> Self {
Self {
info: self.info,
state: self.state.clone(),
track: self.track.clone(),
stats: self.stats.clone(),
}
}
}
impl Drop for Producer {
fn drop(&mut self) {
if !self.state.is_last() {
return;
}
if let Ok(mut state) = modify(&self.state)
&& !state.fin
{
tracing::warn!(
sequence = self.info.sequence,
"group::Producer dropped without finish() or abort()"
);
state.release();
}
}
}
struct Prefetch {
frames: [MaybeUninit<Frame>; Self::CAP],
pos: usize,
len: usize,
}
impl Prefetch {
const CAP: usize = 8;
fn pop(&mut self) -> Option<Frame> {
if self.pos == self.len {
return None;
}
let frame = unsafe { self.frames[self.pos].assume_init_read() };
self.pos += 1;
Some(frame)
}
fn fill(&mut self, frames: impl Iterator<Item = Frame>) {
debug_assert_eq!(self.pos, self.len, "fill on a non-empty batch would leak frames");
self.pos = 0;
self.len = 0;
for frame in frames.take(Self::CAP) {
self.frames[self.len].write(frame);
self.len += 1;
}
}
fn buffered(&self) -> (u64, u64) {
let mut bytes = 0u64;
for slot in &self.frames[self.pos..self.len] {
bytes += unsafe { slot.assume_init_ref() }.payload.len() as u64;
}
((self.len - self.pos) as u64, bytes)
}
}
impl Default for Prefetch {
fn default() -> Self {
Self {
frames: [const { MaybeUninit::uninit() }; Self::CAP],
pos: 0,
len: 0,
}
}
}
impl Drop for Prefetch {
fn drop(&mut self) {
for slot in &mut self.frames[self.pos..self.len] {
unsafe { slot.assume_init_drop() };
}
}
}
pub struct Consumer {
state: kio::Consumer<GroupState>,
info: Info,
track: track::Info,
index: usize,
prefetch: Prefetch,
stats: stats::Meter,
}
impl Clone for Consumer {
fn clone(&self) -> Self {
Self {
state: self.state.clone(),
info: self.info,
track: self.track.clone(),
index: self.index,
prefetch: Prefetch::default(),
stats: self.stats.clone(),
}
}
}
impl std::ops::Deref for Consumer {
type Target = Info;
fn deref(&self) -> &Self::Target {
&self.info
}
}
impl Consumer {
pub(crate) fn with_meter(mut self, meter: stats::Meter) -> Self {
meter.group();
self.stats = meter;
self
}
pub fn timescale(&self) -> Timescale {
self.track.timescale
}
fn poll<F, R>(&self, waiter: &kio::Waiter, f: F) -> Poll<Result<R>>
where
F: Fn(&kio::Ref<'_, GroupState>) -> Poll<Result<R>>,
{
Poll::Ready(match ready!(self.state.poll(waiter, f)) {
Ok(res) => res,
Err(state) => Err(state.abort.clone().unwrap_or(Error::Dropped)),
})
}
pub async fn next_frame(&mut self) -> Result<Option<frame::Consumer>> {
kio::wait(|waiter| self.poll_next_frame(waiter)).await
}
pub fn poll_next_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Consumer>>> {
if let Some(frame) = self.prefetch.pop() {
self.index += 1;
let info = frame::Info {
size: frame.payload.len() as u64,
timestamp: frame.timestamp,
};
let source = frame::Source::Complete(frame.payload);
return Poll::Ready(Ok(Some(frame::Consumer::new(self.state.clone(), info, source))));
}
let index = self.index;
let Some((info, source)) = ready!(self.poll(waiter, |state| state.poll_frame_source(index))?) else {
return Poll::Ready(Ok(None));
};
self.index += 1;
self.stats.frames(1);
Poll::Ready(Ok(Some(
frame::Consumer::new(self.state.clone(), info, source).with_meter(self.stats.clone()),
)))
}
pub fn poll_read_frame(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<frame::Frame>>> {
if let Some(frame) = self.prefetch.pop() {
self.index += 1;
return Poll::Ready(Ok(Some(frame)));
}
let index = self.index;
let prefetch = &mut self.prefetch;
let res = self.state.poll(waiter, |state| {
if index < state.offset {
return Poll::Ready(Err(Error::Lagged));
}
let local = (index - state.offset).min(state.frames.len());
prefetch.fill(state.frames.range(local..).cloned());
if prefetch.len > 0 {
state.charge.touch();
return Poll::Ready(Ok(()));
}
if let Some(err) = &state.abort {
return Poll::Ready(Err(err.clone()));
}
if state.fin {
return Poll::Ready(Ok(()));
}
Poll::Pending
});
match ready!(res) {
Ok(Ok(())) => {}
Ok(Err(err)) => return Poll::Ready(Err(err)),
Err(state) => return Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
}
let (frames, bytes) = self.prefetch.buffered();
self.stats.frames(frames);
self.stats.bytes(bytes);
Poll::Ready(Ok(self.prefetch.pop().inspect(|_| {
self.index += 1;
})))
}
pub async fn read_frame(&mut self) -> Result<Option<frame::Frame>> {
if let Some(frame) = self.prefetch.pop() {
self.index += 1;
return Ok(Some(frame));
}
kio::wait(|waiter| self.poll_read_frame(waiter)).await
}
pub fn poll_finished(&mut self, waiter: &kio::Waiter) -> Poll<Result<u64>> {
self.poll(waiter, |state| state.poll_finished())
}
pub async fn finished(&mut self) -> Result<u64> {
kio::wait(|waiter| self.poll_finished(waiter)).await
}
}
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct Fetch {
pub priority: u8,
}
impl Fetch {
pub fn with_priority(mut self, priority: u8) -> Self {
self.priority = priority;
self
}
}
#[cfg(test)]
mod test {
use super::*;
use bytes::Bytes;
use futures::FutureExt;
#[test]
fn basic_frame_reading() {
let mut producer = Info { sequence: 0 }.produce();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"frame0"))
.unwrap();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"frame1"))
.unwrap();
producer.finish().unwrap();
let mut consumer = producer.consume();
let f0 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(f0.size, 6);
let f1 = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(f1.size, 6);
let end = consumer.next_frame().now_or_never().unwrap().unwrap();
assert!(end.is_none());
}
#[test]
fn read_frame_all_at_once() {
let mut producer = Info { sequence: 0 }.produce();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"hello"))
.unwrap();
producer.finish().unwrap();
let mut consumer = producer.consume();
let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.payload, Bytes::from_static(b"hello"));
}
#[test]
fn read_frame_preserves_timestamp() {
let mut producer = Info { sequence: 0 }.produce();
let timestamp = Timestamp::from_micros(20_000).unwrap();
producer.write_frame(timestamp, Bytes::from_static(b"hello")).unwrap();
producer.finish().unwrap();
let mut consumer = producer.consume();
let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.timestamp.as_micros(), 20_000);
assert_eq!(frame.payload, Bytes::from_static(b"hello"));
}
#[test]
fn chunked_frame_reads_whole() {
let mut producer = Info { sequence: 0 }.produce();
{
let mut frame = producer
.create_frame(frame::Info {
size: 10,
timestamp: Timestamp::ZERO,
})
.unwrap();
frame.write(Bytes::from_static(b"hello")).unwrap();
frame.write(Bytes::from_static(b"world")).unwrap();
frame.finish().unwrap();
}
producer.finish().unwrap();
let mut consumer = producer.consume();
let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.payload, Bytes::from_static(b"helloworld"));
}
#[test]
fn chunked_frame_streams_partial() {
let mut producer = Info { sequence: 0 }.produce();
let mut consumer = producer.consume();
let mut frame = producer
.create_frame(frame::Info {
size: 6,
timestamp: Timestamp::ZERO,
})
.unwrap();
frame.write(Bytes::from_static(b"foo")).unwrap();
let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
let c1 = f.read_chunk().now_or_never().unwrap().unwrap();
assert_eq!(c1, Some(Bytes::from_static(b"foo")));
assert!(f.read_chunk().now_or_never().is_none());
frame.write(Bytes::from_static(b"bar")).unwrap();
frame.finish().unwrap();
let c2 = f.read_chunk().now_or_never().unwrap().unwrap();
assert_eq!(c2, Some(Bytes::from_static(b"bar")));
let c3 = f.read_chunk().now_or_never().unwrap().unwrap();
assert_eq!(c3, None);
}
#[test]
fn group_finish_returns_none() {
let mut producer = Info { sequence: 0 }.produce();
producer.finish().unwrap();
let mut consumer = producer.consume();
let end = consumer.next_frame().now_or_never().unwrap().unwrap();
assert!(end.is_none());
}
#[test]
fn abort_propagates() {
let producer = Info { sequence: 0 }.produce();
let mut consumer = producer.consume();
producer.abort(crate::Error::Cancel).unwrap();
let result = consumer.next_frame().now_or_never().unwrap();
assert!(matches!(result, Err(crate::Error::Cancel)));
}
#[test]
fn abort_clears_cached_frames() {
let mut producer = Info { sequence: 0 }.produce();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
.unwrap();
let _consumer = producer.consume();
assert_eq!(producer.state.read().frames.len(), 1);
producer.clone().abort(crate::Error::Cancel).unwrap();
let state = producer.state.read();
assert!(state.frames.is_empty(), "cached frames should be dropped on abort");
assert_eq!(state.cache, 0);
}
#[test]
fn drop_unfinished_clears_cached_frames() {
let producer = Info { sequence: 0 }.produce();
let mut writer = producer.clone();
writer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
.unwrap();
let mut consumer = producer.consume();
assert_eq!(producer.state.read().frames.len(), 1);
drop(writer);
drop(producer);
let result = consumer.next_frame().now_or_never().unwrap();
assert!(matches!(result, Err(crate::Error::Dropped)));
}
#[test]
fn drop_finished_keeps_cached_frames() {
let mut producer = Info { sequence: 0 }.produce();
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
.unwrap();
producer.finish().unwrap();
let mut consumer = producer.consume();
drop(producer);
let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.payload, Bytes::from_static(b"data"));
}
#[tokio::test]
async fn pending_then_ready() {
let mut producer = Info { sequence: 0 }.produce();
let mut consumer = producer.consume();
assert!(consumer.next_frame().now_or_never().is_none());
producer
.write_frame(Timestamp::ZERO, Bytes::from_static(b"data"))
.unwrap();
producer.finish().unwrap();
let frame = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.size, 4);
}
#[test]
fn eviction_drops_old_frames() {
let mut producer = Info { sequence: 0 }.produce();
let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
producer.write_frame(Timestamp::ZERO, big).unwrap();
let state = producer.state.read();
assert_eq!(state.offset, 1);
assert_eq!(state.frames.len(), 1);
assert_eq!(state.frames[0].payload.len(), MAX_GROUP_CACHE as usize);
}
#[test]
fn next_frame_returns_cache_full_on_tombstone() {
let mut producer = Info { sequence: 0 }.produce();
let big = Bytes::from(vec![0u8; MAX_GROUP_CACHE as usize]);
producer.write_frame(Timestamp::ZERO, big.clone()).unwrap();
producer.write_frame(Timestamp::ZERO, big).unwrap();
let mut consumer = producer.consume();
let result = consumer.next_frame().now_or_never().unwrap();
assert!(matches!(result, Err(crate::Error::Lagged)));
}
#[test]
fn no_eviction_under_budget() {
let mut producer = Info { sequence: 0 }.produce();
for _ in 0..100_000 {
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
}
producer.finish().unwrap();
let state = producer.state.read();
assert_eq!(state.offset, 0);
assert_eq!(state.frames.len(), 100_000);
}
#[test]
fn clone_consumer_independent() {
let mut producer = Info { sequence: 0 }.produce();
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
let mut c1 = producer.consume();
let _ = c1.next_frame().now_or_never().unwrap().unwrap().unwrap();
let mut c2 = c1.clone();
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
producer.finish().unwrap();
let f = c2.next_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(f.size, 1);
let end = c2.next_frame().now_or_never().unwrap().unwrap();
assert!(end.is_none());
}
#[test]
fn read_frame_crosses_prefetch_batches() {
let n = Prefetch::CAP * 3 + 5;
let mut producer = Info { sequence: 0 }.produce();
for i in 0..n {
producer
.write_frame(Timestamp::ZERO, Bytes::from(vec![i as u8; 4]))
.unwrap();
}
producer.finish().unwrap();
let mut consumer = producer.consume();
for i in 0..n {
let frame = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(frame.payload, Bytes::from(vec![i as u8; 4]));
}
assert!(consumer.read_frame().now_or_never().unwrap().unwrap().is_none());
}
#[test]
fn interleave_read_and_next_frame() {
let mut producer = Info { sequence: 0 }.produce();
for i in 0..5u8 {
producer.write_frame(Timestamp::ZERO, Bytes::from(vec![i; 1])).unwrap();
}
producer.finish().unwrap();
let mut consumer = producer.consume();
let f0 = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
assert_eq!(f0.payload, Bytes::from(vec![0u8; 1]));
for i in 1..5u8 {
let mut f = consumer.next_frame().now_or_never().unwrap().unwrap().unwrap();
let data = f.read_all().now_or_never().unwrap().unwrap();
assert_eq!(data, Bytes::from(vec![i; 1]));
}
assert!(consumer.next_frame().now_or_never().unwrap().unwrap().is_none());
}
#[test]
fn read_frame_past_cleared_frames_does_not_panic() {
let mut producer = Info { sequence: 0 }.produce();
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"a")).unwrap();
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"b")).unwrap();
let mut consumer = producer.consume();
consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
producer.abort(Error::Cancel).unwrap();
let result = consumer.read_frame().now_or_never().unwrap();
assert!(matches!(result, Err(Error::Cancel)), "expected Cancel, got {result:?}");
}
#[test]
fn drop_with_partial_batch() {
let mut producer = Info { sequence: 0 }.produce();
for _ in 0..Prefetch::CAP {
producer.write_frame(Timestamp::ZERO, Bytes::from_static(b"x")).unwrap();
}
producer.finish().unwrap();
let mut consumer = producer.consume();
let _ = consumer.read_frame().now_or_never().unwrap().unwrap().unwrap();
drop(consumer);
}
#[test]
fn create_frame_converts_mismatched_scale() {
use crate::{Timescale, Timestamp};
let mut producer = Producer::new(
Info { sequence: 0 },
track::Info::default().with_timescale(Timescale::MICRO),
);
let frame = frame::Info {
size: 3,
timestamp: Timestamp::from_millis(1).unwrap(), };
let writer = producer.create_frame(frame).unwrap();
assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
assert_eq!(writer.timestamp.value(), 1000);
}
#[tokio::test]
async fn create_frame_converts_current_timestamp() {
use crate::Timescale;
let mut producer = Producer::new(
Info { sequence: 0 },
track::Info::default().with_timescale(Timescale::MICRO),
);
let writer = producer
.create_frame(frame::Info {
size: 3,
timestamp: Timestamp::now(),
})
.unwrap();
assert_eq!(writer.timestamp.scale(), Timescale::MICRO);
assert!(!writer.timestamp.is_zero(), "local clock should be non-zero");
}
#[test]
fn create_frame_rejects_oversized() {
let mut producer = Info { sequence: 0 }.produce();
let result = producer.create_frame(frame::Info {
size: MAX_GROUP_CACHE + 1,
timestamp: Timestamp::ZERO,
});
assert!(matches!(result, Err(Error::FrameTooLarge)));
}
}