use std::task::Poll;
use crate::{Error, Result};
#[derive(Default)]
struct State {
bitrate: Option<u64>,
abort: Option<Error>,
}
#[derive(Clone)]
pub struct Producer {
state: kio::Producer<State>,
}
impl Producer {
pub fn new() -> Self {
Self {
state: kio::Producer::default(),
}
}
pub fn set(&self, bitrate: Option<u64>) -> Result<()> {
let mut state = self.modify()?;
if state.bitrate != bitrate {
state.bitrate = bitrate;
}
Ok(())
}
pub fn consume(&self) -> Consumer {
Consumer {
state: self.state.consume(),
last: None,
}
}
pub fn abort(&self, err: Error) -> Result<()> {
let mut state = self.modify()?;
state.abort = Some(err);
state.close();
Ok(())
}
pub async fn closed(&self) {
self.state.closed().await
}
pub async fn unused(&self) -> Result<()> {
kio::wait(|waiter| self.poll_unused(waiter)).await
}
pub fn poll_unused(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
self.state.poll_unused(waiter).map(|used| match used {
Some(()) => Ok(()),
None => Err(self.close_error()),
})
}
pub async fn used(&self) -> Result<()> {
kio::wait(|waiter| self.poll_used(waiter)).await
}
pub fn poll_used(&self, waiter: &kio::Waiter) -> Poll<Result<()>> {
self.state.poll_used(waiter).map(|used| match used {
Some(()) => Ok(()),
None => Err(self.close_error()),
})
}
fn modify(&self) -> Result<kio::Mut<'_, State>> {
self.state
.write()
.map_err(|r| r.abort.clone().unwrap_or(Error::Dropped))
}
fn close_error(&self) -> Error {
self.state.read().abort.clone().unwrap_or(Error::Dropped)
}
}
impl Default for Producer {
fn default() -> Self {
Self::new()
}
}
#[derive(Clone)]
pub struct Consumer {
state: kio::Consumer<State>,
last: Option<u64>,
}
impl Consumer {
pub fn peek(&self) -> Option<u64> {
self.state.read().bitrate
}
pub fn poll_changed(&mut self, waiter: &kio::Waiter) -> Poll<Result<Option<u64>>> {
let last = self.last;
match self.state.poll(waiter, |state| {
if state.bitrate != last {
Poll::Ready(state.bitrate)
} else {
Poll::Pending
}
}) {
Poll::Ready(Ok(bitrate)) => {
self.last = bitrate;
Poll::Ready(Ok(bitrate))
}
Poll::Ready(Err(state)) => Poll::Ready(Err(state.abort.clone().unwrap_or(Error::Dropped))),
Poll::Pending => Poll::Pending,
}
}
pub async fn changed(&mut self) -> Result<Option<u64>> {
kio::wait(|waiter| self.poll_changed(waiter)).await
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn closed_is_distinct_from_unavailable() {
let producer = Producer::new();
let mut consumer = producer.consume();
producer.set(Some(1_000_000)).unwrap();
assert_eq!(consumer.changed().await.unwrap(), Some(1_000_000));
producer.set(None).unwrap();
assert_eq!(consumer.changed().await.unwrap(), None);
producer.abort(Error::Cancel).unwrap();
assert!(consumer.changed().await.is_err());
assert!(consumer.changed().await.is_err());
}
}