#[cfg(not(target_os = "macos"))]
pub(super) use threaded::Sink;
#[cfg(target_os = "macos")]
pub(super) use inline::Sink;
#[cfg(not(target_os = "macos"))]
mod threaded {
use std::thread::JoinHandle;
use bytes::Bytes;
use tokio::sync::{mpsc, oneshot};
use super::super::encoder::{self, Encoder};
use crate::Error;
use crate::frame::Frame;
enum Request {
Encode {
frame: Frame,
keyframe: bool,
resp: oneshot::Sender<Result<Vec<Bytes>, Error>>,
},
SetBitrate {
bitrate: u64,
resp: oneshot::Sender<Result<(), Error>>,
},
}
pub(in crate::encode) struct Sink {
tx: Option<mpsc::UnboundedSender<Request>>,
handle: Option<JoinHandle<()>>,
name: String,
}
impl Sink {
pub(in crate::encode) async fn open(config: &encoder::Config) -> Result<Self, Error> {
let (req_tx, mut req_rx) = mpsc::unbounded_channel::<Request>();
let (ready_tx, ready_rx) = oneshot::channel::<Result<String, Error>>();
let config = config.clone();
let handle = std::thread::spawn(move || {
let mut encoder = match Encoder::new(&config) {
Ok(encoder) => encoder,
Err(err) => {
let _ = ready_tx.send(Err(err));
return;
}
};
if ready_tx.send(Ok(encoder.name().to_string())).is_err() {
return;
}
while let Some(req) = req_rx.blocking_recv() {
match req {
Request::Encode { frame, keyframe, resp } => {
let _ = resp.send(encoder.encode_raw(&frame, keyframe));
}
Request::SetBitrate { bitrate, resp } => {
let _ = resp.send(encoder.set_bitrate(bitrate));
}
}
}
});
match ready_rx.await {
Ok(Ok(name)) => Ok(Self {
tx: Some(req_tx),
handle: Some(handle),
name,
}),
Ok(Err(err)) => Err(err),
Err(_) => {
let _ = handle.join();
Err(Error::Codec(anyhow::anyhow!("encode thread exited before opening")))
}
}
}
pub(in crate::encode) fn name(&self) -> &str {
&self.name
}
pub(in crate::encode) async fn encode(&mut self, frame: Frame, keyframe: bool) -> Result<Vec<Bytes>, Error> {
self.request(|resp| Request::Encode { frame, keyframe, resp }).await
}
pub(in crate::encode) async fn set_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
self.request(|resp| Request::SetBitrate { bitrate, resp }).await
}
async fn request<T>(
&self,
build: impl FnOnce(oneshot::Sender<Result<T, Error>>) -> Request,
) -> Result<T, Error> {
let (resp_tx, resp_rx) = oneshot::channel();
self.tx
.as_ref()
.ok_or_else(gone)?
.send(build(resp_tx))
.map_err(|_| gone())?;
resp_rx.await.map_err(|_| gone())?
}
}
impl Drop for Sink {
fn drop(&mut self) {
self.tx.take();
if let Some(handle) = self.handle.take() {
let _ = handle.join();
}
}
}
fn gone() -> Error {
Error::Codec(anyhow::anyhow!("encode thread stopped unexpectedly"))
}
}
#[cfg(target_os = "macos")]
mod inline {
use bytes::Bytes;
use super::super::encoder::{self, Encoder};
use crate::Error;
use crate::frame::Frame;
pub(in crate::encode) struct Sink(Encoder);
impl Sink {
pub(in crate::encode) async fn open(config: &encoder::Config) -> Result<Self, Error> {
Ok(Self(Encoder::new(config)?))
}
pub(in crate::encode) fn name(&self) -> &str {
self.0.name()
}
pub(in crate::encode) async fn encode(&mut self, frame: Frame, keyframe: bool) -> Result<Vec<Bytes>, Error> {
self.0.encode_raw(&frame, keyframe)
}
pub(in crate::encode) async fn set_bitrate(&mut self, bitrate: u64) -> Result<(), Error> {
self.0.set_bitrate(bitrate)
}
}
}