Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 2 additions & 2 deletions backend.go
Original file line number Diff line number Diff line change
Expand Up @@ -86,8 +86,8 @@ type backendVideoFrame struct {

const maxPooledVideoBuffers = 8

// backendVideoBufferPool keeps decoded RGBA storage independent from the native
// scaler-owned frame, which is overwritten by the next Scale call.
// backendVideoBufferPool reuses the tightly packed RGBA storage filled by the
// native scaler and later uploaded to Ebitengine.
type backendVideoBufferPool struct {
mutex sync.Mutex
buffers [][]byte
Expand Down
50 changes: 20 additions & 30 deletions backend_ffi.go
Original file line number Diff line number Diff line change
Expand Up @@ -41,8 +41,11 @@ type ffiDecoder struct {
decoder *ffmpeg.Decoder
info backendMediaInfo

outputSampleRate int
scaler *ffmpeg.Scaler
outputSampleRate int
scaler *ffmpeg.Scaler
// scaleTarget holds FFmpeg's reference to the most recently filled pooled
// RGBA buffer. WrapBuffer releases that reference before each reuse.
scaleTarget ffmpeg.Frame
scalerWidth int
scalerHeight int
scalerFormat ffmpeg.PixelFormat
Expand Down Expand Up @@ -149,7 +152,6 @@ func (d *ffiDecoder) ReadFrame(ctx context.Context, videoBuffers *backendVideoBu
if frame == nil {
return backendFrame{}, io.EOF
}

converted, ok, err := d.convertFrameLocked(frame, videoBuffers)
if err != nil {
return backendFrame{}, err
Expand Down Expand Up @@ -199,26 +201,20 @@ func (d *ffiDecoder) convertVideoFrameLocked(frame *ffmpeg.FrameWrapper, videoBu
d.scalerFormat = sourceFormat
}

scaled, err := d.scaler.Scale(frame.Raw())
if err != nil {
return backendFrame{}, false, err
}
// Scale returns a scaler-owned raw frame that is overwritten by the next
// call. WrapFrame exposes its data pointer and stride; copyFFmpegRGBA detaches
// it into reusable, tightly packed storage for ebiten.Image.WritePixels.
wrapped := ffmpeg.WrapFrame(scaled, ffmpeg.MediaTypeVideo)
stride := wrapped.Linesize(0)
rowSize := width * 4
if stride < rowSize {
return backendFrame{}, false, fmt.Errorf("avebi: go-ffmpeg-ffi returned RGBA stride %d for row size %d", stride, rowSize)
}
data := wrapped.Data(0)
if len(data) < stride*height {
return backendFrame{}, false, fmt.Errorf("avebi: go-ffmpeg-ffi returned a truncated RGBA frame")
rgba := videoBuffers.get(rowSize * height)
if err := d.scaleTarget.WrapBuffer(rgba, width, height, ffmpeg.PixelFormatRGBA); err != nil {
videoBuffers.put(rgba)
return backendFrame{}, false, fmt.Errorf("wrap RGBA output buffer: %w", err)
}

rgba := videoBuffers.get(rowSize * height)
copyFFmpegRGBA(rgba, data, stride, rowSize, height)
if err := d.scaler.ScaleTo(d.scaleTarget, frame.Raw()); err != nil {
// Release the native reference before returning this buffer to the pool.
_ = d.scaleTarget.Free()
d.scaleTarget = ffmpeg.Frame{}
videoBuffers.put(rgba)
return backendFrame{}, false, err
}

pts := ffmpegPTS(frame.PTS(), d.decoder.VideoStream().TimeBase)
duration := d.info.Video.FrameDuration()
Expand All @@ -235,16 +231,6 @@ func (d *ffiDecoder) convertVideoFrameLocked(frame *ffmpeg.FrameWrapper, videoBu
}, true, nil
}

func copyFFmpegRGBA(dst, src []byte, stride, rowSize, height int) {
if stride == rowSize {
copy(dst, src[:len(dst)])
return
}
for row := 0; row < height; row++ {
copy(dst[row*rowSize:(row+1)*rowSize], src[row*stride:row*stride+rowSize])
}
}

func (d *ffiDecoder) convertAudioFrameLocked(frame *ffmpeg.FrameWrapper) (backendFrame, bool, error) {
if d.info.Audio == nil || d.outputSampleRate <= 0 {
return backendFrame{}, false, nil
Expand Down Expand Up @@ -399,6 +385,10 @@ func (d *ffiDecoder) Close() error {
errs = append(errs, d.scaler.Close())
d.scaler = nil
}
if !d.scaleTarget.IsNil() {
errs = append(errs, d.scaleTarget.Free())
d.scaleTarget = ffmpeg.Frame{}
}
if d.decoder != nil {
errs = append(errs, d.decoder.Close())
d.decoder = nil
Expand Down
17 changes: 17 additions & 0 deletions backend_ffi_integration_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -59,13 +59,18 @@ func TestFFmpegBackendMedia(t *testing.T) {
t.Skip("set AVEBI_TEST_MEDIA to run the FFmpeg integration test")
}

pinnedBefore := ffmpeg.WrappedBufferMemoryUsage()
decoder, err := newMediaBackend().Open(context.Background(), mediaPath, backendOpenOptions{
OutputSampleRate: 48_000,
})
if err != nil {
t.Fatalf("open media: %v", err)
}
closed := false
t.Cleanup(func() {
if closed {
return
}
if err := decoder.Close(); err != nil {
t.Errorf("close media: %v", err)
}
Expand Down Expand Up @@ -93,6 +98,18 @@ func TestFFmpegBackendMedia(t *testing.T) {
t.Fatalf("seek to %s: %v", seekTarget, err)
}
validateFFmpegSeek(t, decoder, seekTarget, wantAudio)

pinnedDuringDecode := ffmpeg.WrappedBufferMemoryUsage()
if pinnedDuringDecode.PinnedBuffers != pinnedBefore.PinnedBuffers+1 {
t.Fatalf("wrapped video buffers while decoding = %d, want %d", pinnedDuringDecode.PinnedBuffers, pinnedBefore.PinnedBuffers+1)
}
if err := decoder.Close(); err != nil {
t.Fatalf("close media: %v", err)
}
closed = true
if pinnedAfterClose := ffmpeg.WrappedBufferMemoryUsage(); pinnedAfterClose != pinnedBefore {
t.Fatalf("wrapped video buffer usage after Close = %+v, want %+v", pinnedAfterClose, pinnedBefore)
}
}

func TestFFmpegPlayerWithDisabledAudioMedia(t *testing.T) {
Expand Down
74 changes: 61 additions & 13 deletions controller_ffi.go
Original file line number Diff line number Diff line change
Expand Up @@ -51,9 +51,13 @@ type playbackController interface {
var _ playbackController = (*ffiLocalController)(nil)

type ffiLocalController struct {
mutex sync.Mutex
decoder mediaDecoder
info backendMediaInfo
// mutex protects playback state. Native decoder calls are serialized
// separately so an audio read cannot hold playback state while decoding.
mutex sync.Mutex
readMutex sync.Mutex
decoderMutex sync.Mutex
decoder mediaDecoder
info backendMediaInfo

state PlaybackState
looping bool
Expand All @@ -79,6 +83,10 @@ type ffiLocalController struct {
muted bool

decodeErr error

// decodeGeneration invalidates a frame decoded across a seek, reset, or
// close while the playback mutex was released.
decodeGeneration uint64
}

func newFFmpegLocalController(decoder mediaDecoder) *ffiLocalController {
Expand Down Expand Up @@ -118,7 +126,7 @@ func (c *ffiLocalController) Play() error {
if err := c.noLockCloseAudioPlayer(); err != nil {
return err
}
if err := c.decoder.Seek(0); err != nil {
if err := c.seekDecoder(0); err != nil {
return err
}
c.noLockResetPlayback(0)
Expand Down Expand Up @@ -172,7 +180,7 @@ func (c *ffiLocalController) Stop() error {
if err := c.noLockCloseAudioPlayer(); err != nil {
return err
}
if err := c.decoder.Seek(0); err != nil {
if err := c.seekDecoder(0); err != nil {
return err
}
c.noLockResetPlayback(0)
Expand All @@ -187,9 +195,10 @@ func (c *ffiLocalController) Close() error {
return nil
}
c.closed = true
c.decodeGeneration++
playerErr := c.noLockCloseAudioPlayer()
c.noLockRecycleVideoFrames()
decoderErr := c.decoder.Close()
decoderErr := c.closeDecoder()
c.videoBuffers.clear()
return errors.Join(playerErr, decoderErr)
}
Expand All @@ -209,7 +218,7 @@ func (c *ffiLocalController) Seek(position time.Duration) (*backendFrame, error)
if err := c.noLockCloseAudioPlayer(); err != nil {
return nil, err
}
if err := c.decoder.Seek(position); err != nil {
if err := c.seekDecoder(position); err != nil {
return nil, err
}
c.noLockResetPlayback(position)
Expand Down Expand Up @@ -259,7 +268,7 @@ func (c *ffiLocalController) Seek(position time.Duration) (*backendFrame, error)
// requested target.
func (c *ffiLocalController) noLockPrimeAfterSeek() (*backendFrame, error) {
for {
frame, err := c.decoder.ReadFrame(context.Background(), &c.videoBuffers)
frame, err := c.readDecoderFrame(context.Background())
if err != nil {
if errors.Is(err, io.EOF) {
return c.lastVideo, nil
Expand Down Expand Up @@ -322,7 +331,7 @@ func (c *ffiLocalController) noLockPosition(now time.Time) (time.Duration, error
if err := c.noLockCloseAudioPlayer(); err != nil {
return 0, err
}
if err := c.decoder.Seek(0); err != nil {
if err := c.seekDecoder(0); err != nil {
return 0, err
}
c.noLockResetPlayback(0)
Expand All @@ -335,7 +344,7 @@ func (c *ffiLocalController) noLockPosition(now time.Time) (time.Duration, error
return 0, nil
}
position %= c.info.Duration
if err := c.decoder.Seek(position); err != nil {
if err := c.seekDecoder(position); err != nil {
return 0, err
}
c.noLockRecycleVideoFrames()
Expand Down Expand Up @@ -438,11 +447,11 @@ func (c *ffiLocalController) CurrentVideoFrame() (*backendFrame, bool, error) {
}

for c.lastVideo == nil || c.lastVideo.PTS+c.frameDuration() < position {
frame, err := c.decoder.ReadFrame(context.Background(), &c.videoBuffers)
frame, err := c.readDecoderFrame(context.Background())
if err != nil {
if errors.Is(err, io.EOF) {
if c.looping {
if err := c.decoder.Seek(0); err != nil {
if err := c.seekDecoder(0); err != nil {
return nil, false, err
}
c.noLockRecycleVideoFrames()
Expand Down Expand Up @@ -517,13 +526,33 @@ func (c *ffiLocalController) Read(buffer []byte) (int, error) {
if len(buffer)&3 != 0 {
buffer = buffer[:len(buffer)&(math.MaxInt-3)]
}
c.readMutex.Lock()
defer c.readMutex.Unlock()

c.mutex.Lock()
defer c.mutex.Unlock()
if c.closed {
return 0, io.EOF
}

served := c.noLockCopyAudio(buffer)
buffer = buffer[served:]
for len(buffer) > 0 {
frame, err := c.decoder.ReadFrame(context.Background(), &c.videoBuffers)
generation := c.decodeGeneration
// Native decoding can take most of a frame interval. Keep playback state
// available to the Ebitengine thread while the audio reader decodes ahead.
c.mutex.Unlock()
frame, err := c.readDecoderFrame(context.Background())
c.mutex.Lock()
if c.closed {
// Close already cleared the pool; let this final buffer become
// unreachable instead of retaining it in a closed controller.
return served, io.EOF
}
if generation != c.decodeGeneration {
recycleBackendFrame(&c.videoBuffers, &frame)
return served, io.EOF
}
if err != nil {
if errors.Is(err, io.EOF) {
// The decoder can reach EOF while Oto still has buffered PCM to
Expand Down Expand Up @@ -558,6 +587,24 @@ func (c *ffiLocalController) Read(buffer []byte) (int, error) {
return served, nil
}

func (c *ffiLocalController) readDecoderFrame(ctx context.Context) (backendFrame, error) {
c.decoderMutex.Lock()
defer c.decoderMutex.Unlock()
return c.decoder.ReadFrame(ctx, &c.videoBuffers)
}

func (c *ffiLocalController) seekDecoder(position time.Duration) error {
c.decoderMutex.Lock()
defer c.decoderMutex.Unlock()
return c.decoder.Seek(position)
}

func (c *ffiLocalController) closeDecoder() error {
c.decoderMutex.Lock()
defer c.decoderMutex.Unlock()
return c.decoder.Close()
}

func (c *ffiLocalController) noLockCopyAudio(buffer []byte) int {
copied := copy(buffer, c.audioQueue)
if copied == len(c.audioQueue) {
Expand Down Expand Up @@ -605,6 +652,7 @@ func (c *ffiLocalController) noLockCloseAudioPlayer() error {
}

func (c *ffiLocalController) noLockResetPlayback(position time.Duration) {
c.decodeGeneration++
c.noLockRecycleVideoFrames()
c.audioQueue = c.audioQueue[:0]
c.firstAudioPTS = position
Expand Down
Loading
Loading