Skip to content
This repository was archived by the owner on Feb 27, 2020. It is now read-only.

Commit e655143

Browse files
author
Wander Lairson Costa
committed
Implement video feature in docker engine
1 parent 4d93412 commit e655143

7 files changed

Lines changed: 165 additions & 5 deletions

File tree

engines/docker/config.go

Lines changed: 10 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -8,8 +8,9 @@ import (
88
)
99

1010
type configType struct {
11-
DockerSocket string `json:"dockerSocket"`
12-
Privileged string `json:"privileged"`
11+
DockerSocket string `json:"dockerSocket"`
12+
Privileged string `json:"privileged"`
13+
EnableDevices bool `json:"enableDevices"`
1314
}
1415

1516
const (
@@ -49,6 +50,13 @@ var configSchema = schematypes.Object{
4950
privilegedNever,
5051
},
5152
},
53+
"enableDevices": schematypes.Boolean{
54+
Title: "Enable host devices",
55+
Description: util.Markdown(`
56+
When true, this enables the support for host devices inside the container,
57+
such as video and sound.
58+
`),
59+
},
5260
},
5361
Required: []string{
5462
"privileged",

engines/docker/engine.go

Lines changed: 28 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -23,6 +23,7 @@ type engine struct {
2323
config configType
2424
networks *network.Pool
2525
imageCache *imagecache.ImageCache
26+
video *videoDeviceManager
2627
}
2728

2829
type engineProvider struct {
@@ -42,6 +43,8 @@ func (p engineProvider) NewEngine(options engines.EngineOptions) (engines.Engine
4243
var c configType
4344
schematypes.MustValidateAndMap(configSchema, options.Config, &c)
4445

46+
debug(fmt.Sprintf("Devices enabled = %v", c.EnableDevices))
47+
4548
if c.DockerSocket == "" {
4649
c.DockerSocket = "unix:///var/run/docker.sock" // default docker socket
4750
}
@@ -54,20 +57,30 @@ func (p engineProvider) NewEngine(options engines.EngineOptions) (engines.Engine
5457

5558
env := options.Environment
5659
monitor := options.Monitor
60+
var video *videoDeviceManager
61+
if c.EnableDevices {
62+
video, err = newVideoDeviceManager()
63+
if err != nil {
64+
return nil, err
65+
}
66+
}
67+
5768
return &engine{
5869
config: c,
5970
docker: client,
6071
Environment: env,
6172
monitor: monitor,
6273
networks: network.NewPool(client, monitor.WithPrefix("network-pool")),
6374
imageCache: imagecache.New(client, env.GarbageCollector, monitor.WithPrefix("image-cache")),
75+
video: video,
6476
}, nil
6577
}
6678

6779
type payloadType struct {
6880
Image interface{} `json:"image"`
6981
Command []string `json:"command"`
7082
Privileged bool `json:"privileged"`
83+
Devices []string `json:"devices"`
7184
}
7285

7386
func (e *engine) PayloadSchema() schematypes.Object {
@@ -79,6 +92,16 @@ func (e *engine) PayloadSchema() schematypes.Object {
7992
Description: "Command to run inside the container.",
8093
Items: schematypes.String{},
8194
},
95+
"devices": schematypes.Array{
96+
Title: "Devices",
97+
Description: "List of host devices required.",
98+
Items: schematypes.StringEnum{
99+
Options: []string{
100+
"video",
101+
},
102+
},
103+
Unique: true,
104+
},
82105
},
83106
Required: []string{
84107
"image",
@@ -107,6 +130,11 @@ func (e *engine) NewSandboxBuilder(options engines.SandboxOptions) (engines.Sand
107130
var p payloadType
108131
schematypes.MustValidateAndMap(e.PayloadSchema(), options.Payload, &p)
109132

133+
if len(p.Devices) > 0 && e.config.EnableDevices == false {
134+
options.TaskContext.LogError(fmt.Sprintf("Task requests device %v, but device support is not enabled", p.Devices))
135+
return nil, runtime.NewMalformedPayloadError("Devices feature is not enabled")
136+
}
137+
110138
// Check if privileged == true is allowed
111139
switch e.config.Privileged {
112140
case privilegedAllow: // Check scope if p.Privileged is true

engines/docker/sandbox.go

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -15,6 +15,7 @@ import (
1515
"github.com/taskcluster/taskcluster-worker/runtime"
1616
"github.com/taskcluster/taskcluster-worker/runtime/atomics"
1717
"github.com/taskcluster/taskcluster-worker/runtime/ioext"
18+
funk "github.com/thoas/go-funk"
1819
)
1920

2021
const dockerEngineKillTimeout = 5 * time.Second
@@ -32,6 +33,7 @@ type sandbox struct {
3233
taskCtx *runtime.TaskContext
3334
networkHandle *network.Handle
3435
imageHandle *imagecache.ImageHandle
36+
videoDev *device
3537
}
3638

3739
func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
@@ -52,6 +54,21 @@ func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
5254
return nil, errors.Wrap(err, "docker.CreateNetwork failed")
5355
}
5456

57+
devices := []docker.Device{}
58+
var dev *device
59+
if funk.InStrings(sb.payload.Devices, "video") {
60+
dev = sb.e.video.claim()
61+
if dev == nil {
62+
return nil, errors.New("No video device available")
63+
}
64+
debug(fmt.Sprintf("Claimed %s", dev.path))
65+
devices = append(devices, docker.Device{
66+
PathOnHost: dev.path,
67+
PathInContainer: dev.path,
68+
CgroupPermissions: "rwm",
69+
})
70+
}
71+
5572
// Create the container
5673
container, err := sb.e.docker.CreateContainer(docker.CreateContainerOptions{
5774
Config: &docker.Config{
@@ -70,6 +87,7 @@ func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
7087
// to the proxies added to proxyMux above..
7188
ExtraHosts: []string{fmt.Sprintf("taskcluster:%s", networkHandle.Gateway())},
7289
Mounts: sb.mounts,
90+
Devices: devices,
7391
},
7492
NetworkingConfig: &docker.NetworkingConfig{
7593
EndpointsConfig: map[string]*docker.EndpointConfig{
@@ -79,6 +97,7 @@ func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
7997
})
8098
if err != nil {
8199
imageHandle.Release()
100+
sb.e.video.release(dev)
82101
return nil, runtime.NewMalformedPayloadError(
83102
"could not create container: " + err.Error())
84103
}
@@ -87,6 +106,7 @@ func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
87106
storage, err := sb.e.Environment.TemporaryStorage.NewFolder()
88107
if err != nil {
89108
imageHandle.Release()
109+
sb.e.video.release(dev)
90110
monitor.ReportError(err, "failed to create temporary folder")
91111
return nil, runtime.ErrFatalInternalError
92112
}
@@ -101,6 +121,7 @@ func newSandbox(sb *sandboxBuilder) (*sandbox, error) {
101121
"containerId": container.ID,
102122
"networkId": networkHandle.NetworkID(),
103123
}),
124+
videoDev: dev,
104125
}
105126

106127
// attach to the container before starting so that we get all the logs
@@ -277,5 +298,9 @@ func (s *sandbox) dispose() error {
277298
if hasErr {
278299
return runtime.ErrNonFatalInternalError
279300
}
301+
302+
if s.videoDev != nil {
303+
s.videoDev.claimed = false
304+
}
280305
return nil
281306
}

engines/docker/video.go

Lines changed: 51 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,51 @@
1+
package dockerengine
2+
3+
import (
4+
"path/filepath"
5+
"regexp"
6+
7+
"github.com/pkg/errors"
8+
funk "github.com/thoas/go-funk"
9+
)
10+
11+
type device struct {
12+
path string
13+
claimed bool
14+
}
15+
16+
type videoDeviceManager struct {
17+
devices []device
18+
}
19+
20+
func newVideoDeviceManager() (*videoDeviceManager, error) {
21+
matches, err := filepath.Glob("/dev/video*")
22+
if err != nil {
23+
return nil, errors.Wrap(err, "Failed to call filepath.Glob function")
24+
}
25+
26+
r := regexp.MustCompile("/dev/video[0-9]+")
27+
matches = funk.FilterString(matches, r.MatchString)
28+
29+
devices := make([]device, len(matches))
30+
for i := range devices {
31+
devices[i].path = matches[i]
32+
}
33+
34+
return &videoDeviceManager{
35+
devices: devices,
36+
}, nil
37+
}
38+
39+
func (d *videoDeviceManager) claim() *device {
40+
for i := range d.devices {
41+
if !d.devices[i].claimed {
42+
return &d.devices[i]
43+
}
44+
}
45+
46+
return nil
47+
}
48+
49+
func (d *videoDeviceManager) release(dev *device) {
50+
dev.claimed = false
51+
}

engines/docker/video_test.go

Lines changed: 36 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,36 @@
1+
// +build dockervideo
2+
3+
package dockerengine
4+
5+
import (
6+
"testing"
7+
8+
"github.com/taskcluster/taskcluster-worker/engines/enginetest"
9+
)
10+
11+
// Image and tag used in test cases below
12+
const (
13+
videoDockerImageName = "alpine:3.6"
14+
)
15+
16+
var videoProvider = &enginetest.EngineProvider{
17+
Engine: "docker",
18+
Config: `{
19+
"privileged": "allow",
20+
"enableDevices": true
21+
}`,
22+
}
23+
24+
func TestVideo(t *testing.T) {
25+
c := enginetest.LoggingTestCase{
26+
EngineProvider: videoProvider,
27+
Target: "/dev/video0",
28+
TargetPayload: `{
29+
"command": ["sh", "-c", "ls /dev/video0"],
30+
"devices": ["video"],
31+
"image": "` + videoDockerImageName + `"
32+
}`,
33+
}
34+
35+
c.Test()
36+
}

engines/enginetest/logging.go

Lines changed: 9 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -71,7 +71,13 @@ func (c *LoggingTestCase) TestSilentTask() {
7171

7272
// Test will run all logging tests
7373
func (c *LoggingTestCase) Test() {
74-
c.TestLogTarget()
75-
c.TestLogTargetWhenFailing()
76-
c.TestSilentTask()
74+
if len(c.TargetPayload) > 0 {
75+
c.TestLogTarget()
76+
}
77+
if len(c.FailingPayload) > 0 {
78+
c.TestLogTargetWhenFailing()
79+
}
80+
if len(c.SilentPayload) > 0 {
81+
c.TestSilentTask()
82+
}
7783
}

vendor/vendor.json

Lines changed: 6 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -677,6 +677,12 @@
677677
"revision": "d4fa08268f573f25df477fedc813f26bfb833761",
678678
"revisionTime": "2016-08-05T02:05:43Z"
679679
},
680+
{
681+
"checksumSHA1": "UHPK5Fi61Zqvf/ctF4s9hwjHFGU=",
682+
"path": "github.com/thoas/go-funk",
683+
"revision": "d2deeb5709c1da54d5da8c76b2f65421f5ff8de4",
684+
"revisionTime": "2018-05-05T20:14:24Z"
685+
},
680686
{
681687
"checksumSHA1": "HN3pLd5cC+QXkX8j8FsCCB3FzSI=",
682688
"path": "github.com/tinylib/msgp/msgp",

0 commit comments

Comments
 (0)