Skip to content

Commit 4d21fb3

Browse files
authored
Fix ImageSink deadlock that wedged the pipeline on image upload failure (#1336)
After looking into a couple of instances of long terminating pods - and collecting goroutine profiles found a bug inside image sink that wedges the pipeline. The ImageSink consumer goroutine exits on the first upload error. after which nobody reads createdImages. Once the buffer fills the gst bus thread on which the element callback is invoked gets blokced on writing to the image channel. That as a results prevents all pipeline message exchanges and causes loop.Quite to stall waiting for the gst bus callback which is blocked forever now. The fix includes 3 changes: a) consumer failure in handleNewImages reports error but keeps draining the channel until it's closed; b) NewImage becomes non blocking send as we can't risk blocking the gst bus on slow io - given default values that would happen only after 60m worth upload queue at which point there is no value in continuing so default case reports an error; c) moving image config validation to start path and failing early if not correct. Test generated by claude.
1 parent a9c3d86 commit 4d21fb3

4 files changed

Lines changed: 196 additions & 13 deletions

File tree

pkg/config/output_image.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -65,6 +65,14 @@ func (p *PipelineConfig) getImageConfig(images *livekit.ImageOutput, upload egre
6565
return nil, err
6666
}
6767

68+
switch images.FilenameSuffix {
69+
case livekit.ImageFileSuffix_IMAGE_SUFFIX_INDEX,
70+
livekit.ImageFileSuffix_IMAGE_SUFFIX_TIMESTAMP,
71+
livekit.ImageFileSuffix_IMAGE_SUFFIX_NONE_OVERWRITE:
72+
default:
73+
return nil, errors.ErrInvalidInput("filename_suffix")
74+
}
75+
6876
sc, err := p.getStorageConfig(upload)
6977
if err != nil {
7078
return nil, err

pkg/errors/errors.go

Lines changed: 4 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -146,6 +146,10 @@ func ErrUploadFailed(location string, err error) error {
146146
return psrpc.NewErrorf(psrpc.InvalidArgument, "%s upload failed: %w", location, err)
147147
}
148148

149+
func ErrUploadQueueFull(sink string, pending int) error {
150+
return psrpc.NewErrorf(psrpc.Internal, "%s upload queue full (%d pending): storage cannot keep up", sink, pending)
151+
}
152+
149153
func ErrParticipantNotFound(identity string) error {
150154
return psrpc.NewErrorf(psrpc.NotFound, "participant %s not found", identity)
151155
}

pkg/pipeline/sink/image.go

Lines changed: 20 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -75,6 +75,9 @@ func newImageSink(
7575
}
7676

7777
maxPendingUploads := (conf.MaxUploadQueue * 60) / int(o.CaptureInterval)
78+
if maxPendingUploads < 1 {
79+
maxPendingUploads = 1
80+
}
7881
return &ImageSink{
7982
base: &base{
8083
bin: imageBin,
@@ -91,19 +94,19 @@ func newImageSink(
9194

9295
func (s *ImageSink) Start() error {
9396
go func() {
94-
var err error
95-
defer func() {
96-
if err != nil {
97-
s.callbacks.OnError(err)
98-
}
99-
s.done.Break()
100-
}()
97+
defer s.done.Break()
10198

99+
// keep draining after a failure: NewImage sends from the pipeline's bus
100+
// thread, and a send with no receiver would block it forever
101+
var failed bool
102102
for update := range s.createdImages {
103-
err = s.handleNewImage(update)
104-
if err != nil {
103+
if failed {
104+
continue
105+
}
106+
if err := s.handleNewImage(update); err != nil {
105107
logger.Errorw("new image handling failed", err)
106-
return
108+
failed = true
109+
s.callbacks.OnError(err)
107110
}
108111
}
109112
}()
@@ -163,12 +166,16 @@ func (s *ImageSink) NewImage(filepath string, ts uint64) error {
163166

164167
filename := filepath[len(s.LocalDir)+1:]
165168

166-
s.createdImages <- &imageUpdate{
169+
// never block: this is called from the pipeline's bus thread
170+
select {
171+
case s.createdImages <- &imageUpdate{
167172
filename: filename,
168173
timestamp: ts,
174+
}:
175+
return nil
176+
default:
177+
return errors.ErrUploadQueueFull("image", cap(s.createdImages))
169178
}
170-
171-
return nil
172179
}
173180

174181
func (s *ImageSink) UploadManifest(filepath string) (string, bool, error) {

pkg/pipeline/sink/image_test.go

Lines changed: 164 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,164 @@
1+
// Copyright 2026 LiveKit, Inc.
2+
//
3+
// Licensed under the Apache License, Version 2.0 (the "License");
4+
// you may not use this file except in compliance with the License.
5+
// You may obtain a copy of the License at
6+
//
7+
// http://www.apache.org/licenses/LICENSE-2.0
8+
//
9+
// Unless required by applicable law or agreed to in writing, software
10+
// distributed under the License is distributed on an "AS IS" BASIS,
11+
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12+
// See the License for the specific language governing permissions and
13+
// limitations under the License.
14+
15+
package sink
16+
17+
import (
18+
"fmt"
19+
"os"
20+
"path"
21+
"testing"
22+
"time"
23+
24+
"github.com/stretchr/testify/require"
25+
26+
"github.com/livekit/egress/pkg/config"
27+
"github.com/livekit/egress/pkg/gstreamer"
28+
"github.com/livekit/egress/pkg/pipeline/sink/uploader"
29+
"github.com/livekit/protocol/livekit"
30+
)
31+
32+
func newTestImageSink(t *testing.T, capacity int) (*ImageSink, chan error) {
33+
u, err := uploader.New(nil, nil, nil, nil, nil)
34+
require.NoError(t, err)
35+
36+
errCh := make(chan error, 10)
37+
callbacks := &gstreamer.Callbacks{}
38+
callbacks.SetOnError(func(err error) { errCh <- err })
39+
40+
return &ImageSink{
41+
Uploader: u,
42+
ImageConfig: &config.ImageConfig{
43+
ImagesInfo: &livekit.ImagesInfo{},
44+
LocalDir: t.TempDir(),
45+
StorageDir: t.TempDir(),
46+
ImagePrefix: "img",
47+
ImageSuffix: livekit.ImageFileSuffix_IMAGE_SUFFIX_INDEX,
48+
},
49+
conf: &config.PipelineConfig{},
50+
callbacks: callbacks,
51+
createdImages: make(chan *imageUpdate, capacity),
52+
}, errCh
53+
}
54+
55+
func requireNewImage(t *testing.T, s *ImageSink, filepath string) {
56+
done := make(chan error, 1)
57+
go func() {
58+
done <- s.NewImage(filepath, uint64(time.Now().UnixNano()))
59+
}()
60+
select {
61+
case err := <-done:
62+
require.NoError(t, err)
63+
case <-time.After(3 * time.Second):
64+
t.Fatal("NewImage blocked")
65+
}
66+
}
67+
68+
func requireClose(t *testing.T, s *ImageSink) {
69+
closed := make(chan struct{})
70+
go func() {
71+
require.NoError(t, s.Close())
72+
close(closed)
73+
}()
74+
select {
75+
case <-closed:
76+
case <-time.After(3 * time.Second):
77+
t.Fatal("Close blocked")
78+
}
79+
}
80+
81+
// An upload failure must fail the egress exactly once, and the consumer must
82+
// keep accepting images so the producer (the pipeline's bus thread) never blocks.
83+
func TestImageSinkDrainsAfterUploadFailure(t *testing.T) {
84+
s, errCh := newTestImageSink(t, 2)
85+
require.NoError(t, s.Start())
86+
87+
// no local file exists, so handleNewImage fails on upload
88+
requireNewImage(t, s, path.Join(s.LocalDir, "img_00001.jpg"))
89+
90+
select {
91+
case err := <-errCh:
92+
require.Error(t, err)
93+
case <-time.After(5 * time.Second):
94+
t.Fatal("OnError not called after upload failure")
95+
}
96+
97+
// the consumer must keep draining well past channel capacity; sends may
98+
// transiently report a full queue but must never block
99+
for i := 2; i <= 12; i++ {
100+
fp := path.Join(s.LocalDir, fmt.Sprintf("img_%05d.jpg", i))
101+
require.Eventually(t, func() bool {
102+
err := s.NewImage(fp, uint64(time.Now().UnixNano()))
103+
if err != nil {
104+
require.ErrorContains(t, err, "upload queue full")
105+
return false
106+
}
107+
return true
108+
}, 3*time.Second, 10*time.Millisecond, "image %d never accepted", i)
109+
}
110+
111+
requireClose(t, s)
112+
113+
select {
114+
case err := <-errCh:
115+
t.Fatalf("OnError called more than once: %v", err)
116+
default:
117+
}
118+
}
119+
120+
func TestImageSinkNewImageQueueFull(t *testing.T) {
121+
s, _ := newTestImageSink(t, 1)
122+
// consumer not started: first send fills the buffer, second must fail fast
123+
124+
require.NoError(t, s.NewImage(path.Join(s.LocalDir, "img_00001.jpg"), 0))
125+
126+
err := s.NewImage(path.Join(s.LocalDir, "img_00002.jpg"), 0)
127+
require.Error(t, err)
128+
require.ErrorContains(t, err, "upload queue full")
129+
}
130+
131+
func TestImageSinkCloseDrainsPendingUploads(t *testing.T) {
132+
// the default local storage backend resolves storage paths against the
133+
// working directory at construction time
134+
t.Chdir(t.TempDir())
135+
136+
s, errCh := newTestImageSink(t, 4)
137+
s.StorageDir = "images-out"
138+
139+
filenames := []string{"img_00001.jpg", "img_00002.jpg"}
140+
for _, name := range filenames {
141+
require.NoError(t, os.WriteFile(path.Join(s.LocalDir, name), []byte("jpeg"), 0o644))
142+
}
143+
144+
require.NoError(t, s.Start())
145+
for _, name := range filenames {
146+
requireNewImage(t, s, path.Join(s.LocalDir, name))
147+
}
148+
149+
requireClose(t, s)
150+
151+
require.EqualValues(t, len(filenames), s.ImagesInfo.ImageCount)
152+
for _, name := range filenames {
153+
_, err := os.Stat(path.Join(s.StorageDir, name))
154+
require.NoError(t, err, "image not uploaded to storage")
155+
_, err = os.Stat(path.Join(s.LocalDir, name))
156+
require.True(t, os.IsNotExist(err), "local image not removed after upload")
157+
}
158+
159+
select {
160+
case err := <-errCh:
161+
t.Fatalf("unexpected error: %v", err)
162+
default:
163+
}
164+
}

0 commit comments

Comments
 (0)