-
Notifications
You must be signed in to change notification settings - Fork 0
Expand file tree
/
Copy pathdecoder_state.go
More file actions
158 lines (137 loc) · 3.31 KB
/
Copy pathdecoder_state.go
File metadata and controls
158 lines (137 loc) · 3.31 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
//go:build amd64 || arm64
package ffmpeg
import (
"errors"
"github.com/bstkhq/go-ffmpeg-ffi/avcodec"
"github.com/bstkhq/go-ffmpeg-ffi/avutil"
)
var errDecoderProtocolStalled = errors.New("ffmpeg: decoder send/receive protocol stalled")
// decoderCodecState owns packets until FFmpeg accepts them and tracks the
// one-way transition from regular input to draining at end of stream.
type decoderCodecState struct {
pending []avcodec.Packet
flushRequested bool
flushSent bool
drained bool
sendPacket func(avcodec.Context, avcodec.Packet) error
receiveFrame func(avcodec.Context, avutil.Frame) error
unrefFrame func(avutil.Frame)
freePacket func(*avcodec.Packet)
}
func (s *decoderCodecState) enqueueOwned(packet avcodec.Packet) error {
if packet == nil {
return nil
}
if s.flushRequested || s.flushSent || s.drained {
return errors.New("ffmpeg: cannot submit a packet after decoder flush")
}
s.pending = append(s.pending, packet)
return nil
}
func (s *decoderCodecState) requestFlush() {
s.flushRequested = true
}
func (s *decoderCodecState) next(ctx avcodec.Context, frame avutil.Frame) (bool, error) {
if s.drained {
return false, nil
}
for {
s.unref(frame)
err := s.receive(ctx, frame)
if err == nil {
return true, nil
}
if avutil.IsEOF(err) {
s.drained = true
s.clearPending()
return false, nil
}
if !avutil.IsAgain(err) {
return false, err
}
if len(s.pending) > 0 {
err = s.send(ctx, s.pending[0])
switch {
case err == nil:
s.releaseFirst()
continue
case avutil.IsAgain(err):
return false, errDecoderProtocolStalled
case avutil.IsEOF(err):
s.drained = true
s.clearPending()
return false, err
default:
return false, err
}
}
if !s.flushRequested {
return false, nil
}
if s.flushSent {
return false, errDecoderProtocolStalled
}
err = s.send(ctx, nil)
switch {
case err == nil:
s.flushSent = true
continue
case avutil.IsAgain(err):
return false, errDecoderProtocolStalled
case avutil.IsEOF(err):
s.flushSent = true
s.drained = true
return false, nil
default:
return false, err
}
}
}
func (s *decoderCodecState) reset() {
s.clearPending()
s.flushRequested = false
s.flushSent = false
s.drained = false
}
func (s *decoderCodecState) hasPending() bool {
return len(s.pending) > 0
}
func (s *decoderCodecState) releaseFirst() {
packet := s.pending[0]
copy(s.pending, s.pending[1:])
s.pending[len(s.pending)-1] = nil
s.pending = s.pending[:len(s.pending)-1]
s.free(&packet)
}
func (s *decoderCodecState) clearPending() {
for i := range s.pending {
s.free(&s.pending[i])
}
s.pending = nil
}
func (s *decoderCodecState) send(ctx avcodec.Context, packet avcodec.Packet) error {
if s.sendPacket != nil {
return s.sendPacket(ctx, packet)
}
return avcodec.SendPacket(ctx, packet)
}
func (s *decoderCodecState) receive(ctx avcodec.Context, frame avutil.Frame) error {
if s.receiveFrame != nil {
return s.receiveFrame(ctx, frame)
}
return avcodec.ReceiveFrame(ctx, frame)
}
func (s *decoderCodecState) unref(frame avutil.Frame) {
if s.unrefFrame != nil {
s.unrefFrame(frame)
return
}
avutil.FrameUnref(frame)
}
func (s *decoderCodecState) free(packet *avcodec.Packet) {
if s.freePacket != nil {
s.freePacket(packet)
return
}
avcodec.PacketFree(packet)
}