use std::sync::{Arc, Condvar, Mutex};
use std::time::Duration;
use bytes::Bytes;
use moq_net::Timestamp;
use ndk::hardware_buffer::HardwareBufferUsage;
use ndk::media::image_reader::{AcquireResult, Image, ImageFormat, ImageReader};
use ndk::media::media_codec::{self, DequeuedInputBufferResult, DequeuedOutputBufferInfoResult, MediaCodecDirection};
use ndk::media::media_format::MediaFormat;
use ndk::media_error::MediaError;
use super::{Backend, Codec, Config};
use crate::frame::{
Surface,
android::{HardwareBuffer, Reader},
};
use crate::{Error, Frame};
pub(crate) const NAME: &str = "mediacodec";
const MIME_H264: &str = "video/avc";
const MIME_H265: &str = "video/hevc";
const MIME_AV1: &str = "video/av01";
const KEY_MIME: &str = "mime";
const KEY_WIDTH: &str = "width";
const KEY_HEIGHT: &str = "height";
const KEY_ALLOW_FRAME_DROP: &str = "allow-frame-drop";
const KEY_LOW_LATENCY: &str = "low-latency";
const KEY_PRIORITY: &str = "priority";
const PRIORITY_REALTIME: i32 = 0;
const FLAG_CODEC_CONFIG: u32 = 2;
const FLAG_END_OF_STREAM: u32 = 4;
const DEFAULT_SIZE: (i32, i32) = (1920, 1080);
const QUEUE_DEPTH: i32 = 8;
const INPUT_TIMEOUT: Duration = Duration::from_millis(10);
const OUTPUT_TIMEOUT: Duration = Duration::ZERO;
const DRAIN_TIMEOUT: Duration = Duration::from_millis(50);
const DRAIN_ROUNDS: u32 = 20;
const SUBMIT_ROUNDS: u32 = 20;
pub(crate) struct MediaCodec {
codec: media_codec::MediaCodec,
reader: Arc<Reader>,
images: Arc<ImageSignal>,
image_generation: u64,
pending_images: usize,
fed: bool,
stalled: bool,
}
#[derive(Default)]
struct ImageSignal {
generation: Mutex<u64>,
ready: Condvar,
}
impl ImageSignal {
fn notify(&self) {
let mut generation = self.generation.lock().unwrap_or_else(|e| e.into_inner());
*generation = generation.wrapping_add(1);
self.ready.notify_one();
}
fn current(&self) -> u64 {
*self.generation.lock().unwrap_or_else(|e| e.into_inner())
}
fn wait(&self, generation: u64, timeout: Duration) -> u64 {
let current = self.generation.lock().unwrap_or_else(|e| e.into_inner());
let (current, _) = self
.ready
.wait_timeout_while(current, timeout, |current| *current == generation)
.unwrap_or_else(|e| e.into_inner());
*current
}
}
unsafe impl Send for MediaCodec {}
impl MediaCodec {
pub(crate) fn open(codec: Codec, _config: &Config) -> Result<Box<dyn Backend>, Error> {
let mime = match codec {
Codec::H264 => MIME_H264,
Codec::H265 => MIME_H265,
Codec::Av1 => MIME_AV1,
};
let (width, height) = DEFAULT_SIZE;
let mut reader = ImageReader::new_with_usage(
width,
height,
ImageFormat::YUV_420_888,
HardwareBufferUsage::GPU_SAMPLED_IMAGE | HardwareBufferUsage::CPU_READ_OFTEN,
QUEUE_DEPTH,
)
.map_err(|e| reader_err("create an ImageReader", e))?;
let images = Arc::new(ImageSignal::default());
let ready = images.clone();
reader
.set_image_listener(Box::new(move |_| ready.notify()))
.map_err(|e| reader_err("register the image listener", e))?;
reader
.set_buffer_removed_listener(Box::new(|_, _| {}))
.map_err(|e| reader_err("register the buffer removal listener", e))?;
let surface = reader.window().map_err(|e| reader_err("get the reader's window", e))?;
let decoder = media_codec::MediaCodec::from_decoder_type(mime)
.ok_or_else(|| Error::Codec(anyhow::anyhow!("no MediaCodec decoder for {mime}")))?;
let format = decoder_format(mime, width, height);
decoder
.configure(&format, Some(&surface), MediaCodecDirection::Decoder)
.map_err(|e| codec_err("configure", e))?;
decoder.start().map_err(|e| codec_err("start", e))?;
tracing::info!(decoder = NAME, codec = codec.label(), "opened video decoder");
Ok(Box::new(Self {
codec: decoder,
reader: Arc::new(Reader::new(reader)),
images,
image_generation: 0,
pending_images: 0,
fed: false,
stalled: false,
}))
}
fn submit(&mut self, access_unit: &[u8], timestamp: Timestamp, out: &mut Vec<Frame>) -> Result<(), Error> {
let time = timestamp.as_micros().min(u64::MAX as u128) as u64;
for _ in 0..SUBMIT_ROUNDS {
let queued = match self
.codec
.dequeue_input_buffer(INPUT_TIMEOUT)
.map_err(|e| codec_err("dequeue an input buffer", e))?
{
DequeuedInputBufferResult::Buffer(mut buffer) => {
let target = buffer.buffer_mut();
if target.len() < access_unit.len() {
return Err(Error::Codec(anyhow::anyhow!(
"MediaCodec input buffer is {} bytes, needs {} for this access unit",
target.len(),
access_unit.len()
)));
}
unsafe {
std::ptr::copy_nonoverlapping(
access_unit.as_ptr(),
target.as_mut_ptr().cast::<u8>(),
access_unit.len(),
);
}
self.codec
.queue_input_buffer(buffer, 0, access_unit.len(), time, 0)
.map_err(|e| codec_err("queue an input buffer", e))?;
self.fed = true;
true
}
DequeuedInputBufferResult::TryAgainLater => false,
};
if queued {
return Ok(());
}
self.drain(OUTPUT_TIMEOUT, out)?;
}
Err(Error::Codec(anyhow::anyhow!(
"MediaCodec never freed an input buffer for an access unit"
)))
}
fn drain(&mut self, timeout: Duration, out: &mut Vec<Frame>) -> Result<bool, Error> {
loop {
let ended = match self
.codec
.dequeue_output_buffer(timeout)
.map_err(|e| codec_err("dequeue an output buffer", e))?
{
DequeuedOutputBufferInfoResult::Buffer(buffer) => {
let info = *buffer.info();
let render = info.flags() & FLAG_CODEC_CONFIG == 0 && info.size() > 0;
self.codec
.release_output_buffer(buffer, render)
.map_err(|e| codec_err("release an output buffer", e))?;
if render {
self.pending_images += 1;
}
info.flags() & FLAG_END_OF_STREAM != 0
}
DequeuedOutputBufferInfoResult::TryAgainLater => {
self.collect(out)?;
return Ok(false);
}
DequeuedOutputBufferInfoResult::OutputFormatChanged
| DequeuedOutputBufferInfoResult::OutputBuffersChanged => false,
};
self.collect(out)?;
if ended {
return Ok(true);
}
}
}
fn collect(&mut self, out: &mut Vec<Frame>) -> Result<(), Error> {
loop {
match self
.reader
.acquire_next_image()
.map_err(|e| reader_err("acquire an image", e))?
{
AcquireResult::Image(image) => {
self.pending_images = self.pending_images.saturating_sub(1);
let nanos = image
.timestamp()
.map_err(|e| reader_err("read an image's timestamp", e))?;
let timestamp = Timestamp::from_nanos(nanos.max(0) as u64)?;
let (left, top, width, height) = image_size(&image)?;
let buffer = HardwareBuffer::new(self.reader.clone(), image, left, top, width, height);
out.push(Frame::new(Surface::HardwareBuffer(buffer), timestamp));
}
AcquireResult::NoBufferAvailable => {
self.image_generation = self.images.current();
self.stalled = false;
break;
}
AcquireResult::MaxImagesAcquired => {
if !self.stalled {
self.stalled = true;
tracing::warn!(
decoder = NAME,
depth = QUEUE_DEPTH,
"every decoded frame is still held by a consumer; decoding is stalled"
);
}
break;
}
}
}
Ok(())
}
fn drain_tail(&mut self) -> Result<Vec<Frame>, Error> {
let mut out = Vec::new();
if !self.fed {
return Ok(out);
}
self.signal_end_of_input(&mut out)?;
let mut ended = false;
for _ in 0..DRAIN_ROUNDS {
if self.drain(DRAIN_TIMEOUT, &mut out)? {
ended = true;
break;
}
}
if !ended {
return Err(Error::Codec(anyhow::anyhow!(
"MediaCodec did not reach end of stream within {:?}",
DRAIN_TIMEOUT * DRAIN_ROUNDS
)));
}
for _ in 0..DRAIN_ROUNDS {
self.collect(&mut out)?;
if self.pending_images == 0 {
break;
}
self.image_generation = self.images.wait(self.image_generation, DRAIN_TIMEOUT);
}
if self.pending_images > 0 {
tracing::warn!(
decoder = NAME,
dropped = self.pending_images,
"rendered decoder outputs never arrived at the ImageReader"
);
self.pending_images = 0;
}
Ok(out)
}
fn signal_end_of_input(&mut self, out: &mut Vec<Frame>) -> Result<(), Error> {
for _ in 0..DRAIN_ROUNDS {
let queued = match self
.codec
.dequeue_input_buffer(INPUT_TIMEOUT)
.map_err(|e| codec_err("dequeue an input buffer", e))?
{
DequeuedInputBufferResult::Buffer(buffer) => {
self.codec
.queue_input_buffer(buffer, 0, 0, 0, FLAG_END_OF_STREAM)
.map_err(|e| codec_err("queue end of stream", e))?;
true
}
DequeuedInputBufferResult::TryAgainLater => false,
};
if queued {
return Ok(());
}
self.drain(OUTPUT_TIMEOUT, out)?;
}
Err(Error::Codec(anyhow::anyhow!(
"MediaCodec never freed an input buffer for the end of stream"
)))
}
}
impl Backend for MediaCodec {
fn decode(&mut self, access_unit: Bytes, timestamp: Timestamp, _keyframe: bool) -> Result<Vec<Frame>, Error> {
let mut out = Vec::new();
self.drain(OUTPUT_TIMEOUT, &mut out)?;
self.submit(&access_unit, timestamp, &mut out)?;
self.drain(OUTPUT_TIMEOUT, &mut out)?;
Ok(out)
}
fn flush(&mut self) -> Result<Vec<Frame>, Error> {
let out = self.drain_tail()?;
if self.fed {
self.codec.flush().map_err(|e| codec_err("flush", e))?;
self.fed = false;
}
Ok(out)
}
fn name(&self) -> &str {
NAME
}
}
fn decoder_format(mime: &str, width: i32, height: i32) -> MediaFormat {
let mut format = MediaFormat::new();
format.set_str(KEY_MIME, mime);
format.set_i32(KEY_WIDTH, width);
format.set_i32(KEY_HEIGHT, height);
format.set_i32(KEY_ALLOW_FRAME_DROP, 0);
format.set_i32(KEY_LOW_LATENCY, 1);
format.set_i32(KEY_PRIORITY, PRIORITY_REALTIME);
format
}
fn image_size(image: &Image) -> Result<(u32, u32, u32, u32), Error> {
let width = image.width().map_err(|e| reader_err("read an image's width", e))?;
let height = image.height().map_err(|e| reader_err("read an image's height", e))?;
let crop = image.crop_rect().map_err(|e| reader_err("read an image's crop", e))?;
let (mut left, mut top, mut w, mut h) = (0, 0, width, height);
let (cropped_w, cropped_h) = (crop.right - crop.left, crop.bottom - crop.top);
if crop.left >= 0 && crop.top >= 0 && crop.right <= width && crop.bottom <= height && cropped_w > 0 && cropped_h > 0
{
(left, top, w, h) = (crop.left, crop.top, cropped_w, cropped_h);
}
let (w, h) = (w.max(0) as u32 & !1, h.max(0) as u32 & !1);
if w == 0 || h == 0 {
return Err(Error::Codec(anyhow::anyhow!(
"MediaCodec produced a {width}x{height} image, which is not a picture"
)));
}
Ok((left as u32, top as u32, w, h))
}
fn codec_err(what: &str, error: MediaError) -> Error {
Error::Codec(anyhow::anyhow!("failed to {what} on the MediaCodec decoder: {error}"))
}
fn reader_err(what: &str, error: MediaError) -> Error {
Error::Codec(anyhow::anyhow!(
"failed to {what} on the decoder's ImageReader: {error}"
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
#[ignore = "needs an Android device with MediaCodec H.264 hardware"]
fn decodes_what_the_encoder_produced() {
use crate::encode::{Config as EncodeConfig, Kind as EncodeKind};
let size = crate::Size::new(320, 240);
let mut config = EncodeConfig::new(size.width, size.height, 30);
config.kind = EncodeKind::Named(NAME.to_owned());
let mut encoder = crate::encode::Encoder::new(&config).expect("a MediaCodec encoder");
let i420 = crate::I420::new(
size.width,
size.height,
vec![0x80; crate::I420::len(size.width, size.height)],
)
.unwrap();
let mut decoder = MediaCodec::open(Codec::H264, &Config::new()).expect("a MediaCodec decoder");
let mut frames = Vec::new();
for index in 0..30u64 {
let timestamp = Timestamp::from_micros(index * 33_333).unwrap();
let frame = Frame::new(Surface::I420(i420.clone()), timestamp);
encoder.keyframe();
for encoded in encoder.encode(&frame).unwrap() {
frames.extend(decoder.decode(encoded.payload, encoded.timestamp, true).unwrap());
}
}
for encoded in encoder.flush().unwrap() {
frames.extend(decoder.decode(encoded.payload, encoded.timestamp, true).unwrap());
}
frames.extend(decoder.flush().unwrap());
let timestamp = Timestamp::from_micros(1_000_000).unwrap();
let frame = Frame::new(Surface::I420(i420.clone()), timestamp);
encoder.keyframe();
for encoded in encoder.encode(&frame).unwrap() {
frames.extend(decoder.decode(encoded.payload, encoded.timestamp, true).unwrap());
}
for encoded in encoder.flush().unwrap() {
frames.extend(decoder.decode(encoded.payload, encoded.timestamp, true).unwrap());
}
frames.extend(decoder.flush().unwrap());
let frame = frames.first().expect("at least one decoded frame");
assert_eq!(frame.size(), size);
assert!(
matches!(frame.surface, Surface::HardwareBuffer(_)),
"a decoded picture should stay in its hardware buffer",
);
let frames = frames.into_iter().next().expect("checked above");
let i420 = frames.surface.into_i420().expect("read back to I420");
assert_eq!(i420.len(), crate::I420::len(size.width, size.height));
let luma = &i420[..(size.width * size.height) as usize];
assert!(
luma.iter().all(|byte| byte.abs_diff(0x80) <= 8),
"the read-back luma plane should still be mid-gray",
);
}
}