fix: stop the job's docker socket from becoming a directory (#1215)

Fixes https://gitea.com/gitea/runner/issues/1213

Fix the DooD regression that mounts `/var/run/docker.sock` as a directory. Keep the Docker proxy available through job and post steps. Clean stale resources before opening it, then remove containers before their networks and volumes during teardown.

Use a unique filesystem probe and preserve socket ownership. Fall back to direct access when proxying is unsupported. Preserve exec output and clean up active streams and failed starts.

Add a real Docker job test for mounted socket access, post steps and resource cleanup.

---------

Co-authored-by: silverwind <me@silverwind.io>
Reviewed-on: https://gitea.com/gitea/runner/pulls/1215
Reviewed-by: silverwind <2021+silverwind@noreply.gitea.com>
Co-authored-by: Zettat123 <zettat123@gmail.com>
This commit is contained in:
Zettat123
2026-09-08 04:16:02 +00:00
committed by bircni
co-authored by silverwind
parent ba4d3c5b4f
commit ff9965e940
18 changed files with 935 additions and 295 deletions
+7 -4
View File
@@ -9,6 +9,7 @@ import (
"errors"
"fmt"
"io"
"sync"
"gitea.com/gitea/runner/act/common"
@@ -77,13 +78,15 @@ var ErrContainerNotFound = errors.New("does not exist")
// DockerProxy is a job's docker socket, fronting the daemon's for the job's lifetime.
type DockerProxy struct {
Socket string
close func(context.Context) error
Socket string
close func(context.Context) error
closeOnce sync.Once
closeErr error
}
// Close removes what the job created through the socket, then stops serving it.
func (p *DockerProxy) Close(ctx context.Context) error {
return p.close(ctx)
p.closeOnce.Do(func() { p.closeErr = p.close(ctx) })
return p.closeErr
}
// Info is a snapshot of a container, as of one inspect.
+229 -110
View File
@@ -6,20 +6,22 @@
package container
import (
"bufio"
"bytes"
"context"
"crypto/sha256"
"encoding/hex"
"encoding/json"
"errors"
"fmt"
"io"
"maps"
"mime"
"net"
"net/http"
"net/http/httputil"
"os"
"path/filepath"
"regexp"
"runtime"
"strings"
"sync"
"time"
@@ -35,64 +37,84 @@ import (
const (
jobLabel = "com.gitea.runner.job"
maxCreateBody = 8 << 20
dockerProxyProbeTimeout = 5 * time.Second
)
var (
createPath = regexp.MustCompile(`^(/v[0-9.]+)?/(containers|networks|volumes)/create$`)
rawStreamPath = regexp.MustCompile(`^(/v[0-9.]+)?/(containers/[^/]+/attach|exec/[^/]+/start)$`)
proxyProbe struct {
sync.Mutex
decided bool
dir string
}
)
// DockerProxyDir returns where job proxy sockets live, "" while undecided or when the daemon cannot open the runner's files.
func DockerProxyDir(ctx context.Context) string {
proxyProbe.Lock()
defer proxyProbe.Unlock()
if proxyProbe.decided {
return proxyProbe.dir
func NewDockerProxy(ctx context.Context, job string) *DockerProxy {
if host := os.Getenv("DOCKER_HOST"); runtime.GOOS != "linux" || host != "" && !strings.HasPrefix(host, "unix://") {
return nil
}
dir := filepath.Join(os.TempDir(), "gitea-runner-docker")
ok, err := daemonSeesDir(ctx, dir)
probeCtx, cancel := context.WithTimeout(ctx, dockerProxyProbeTimeout)
defer cancel()
cli, err := GetDockerClient(probeCtx)
if err != nil {
common.Logger(ctx).Debugf("docker proxy probe postponed: %v", err)
return ""
}
proxyProbe.decided = true
if ok {
proxyProbe.dir = dir
} else {
common.Logger(ctx).Infof("the docker daemon cannot reach the runner's filesystem, jobs get the daemon socket directly")
}
return proxyProbe.dir
}
func daemonSeesDir(ctx context.Context, dir string) (bool, error) {
if err := os.MkdirAll(dir, 0o700); err != nil {
return false, err
}
probe := filepath.Join(dir, "probe")
if err := os.WriteFile(probe, nil, 0o600); err != nil {
return false, err
}
cli, err := GetDockerClient(ctx)
if err != nil {
return false, err
return nil
}
defer cli.Close()
daemonSocket, ok := strings.CutPrefix(cli.DaemonHost(), "unix://")
if !ok {
return nil
}
if info, err := os.Stat(daemonSocket); err != nil || info.Mode()&os.ModeSocket == 0 {
return nil
}
dir, err := filepath.Abs(os.TempDir())
if err != nil {
common.Logger(ctx).Infof("docker proxy probe failed, jobs get the daemon socket directly: %v", err)
return nil
}
seen, err := daemonSeesDir(probeCtx, cli, dir)
if err != nil {
common.Logger(ctx).Infof("docker proxy probe failed, jobs get the daemon socket directly: %v", err)
return nil
}
if !seen {
common.Logger(ctx).Infof("the docker daemon cannot reach the runner's temporary filesystem, jobs get the daemon socket directly")
return nil
}
if ctx.Err() != nil {
return nil
}
proxy, err := StartDockerProxy(daemonSocket, dir, job)
if err != nil {
common.Logger(ctx).Warnf("docker proxy not started, the job gets the daemon socket directly: %v", err)
}
return proxy
}
// daemonSeesDir reports whether the daemon opens the files the runner writes in dir,
// which is what a job's proxy socket mounted from there needs.
func daemonSeesDir(ctx context.Context, cli client.APIClient, dir string) (bool, error) {
marker, err := os.CreateTemp(dir, "gitea-runner-probe-")
if err != nil {
return false, err
}
defer func() {
if err := os.Remove(marker.Name()); err != nil {
common.Logger(ctx).Warnf("removing the docker proxy probe marker failed: %v", err)
}
}()
if err := marker.Close(); err != nil {
return false, err
}
images, err := cli.ImageList(ctx, client.ImageListOptions{})
if err != nil {
return false, err
}
if len(images.Items) == 0 {
return false, errors.New("no image to probe with yet")
return false, errors.New("no image available for the docker proxy probe")
}
// creating validates that a bind source exists on the daemon's side, nothing is started
created, err := cli.ContainerCreate(ctx, client.ContainerCreateOptions{
Config: &container.Config{Image: images.Items[0].ID, Cmd: []string{"true"}},
HostConfig: &container.HostConfig{Mounts: []mount.Mount{{Type: mount.TypeBind, Source: probe, Target: "/gitea-runner-probe"}}},
Config: &container.Config{Image: images.Items[0].ID, Cmd: []string{"true"}},
HostConfig: &container.HostConfig{Mounts: []mount.Mount{
{Type: mount.TypeBind, Source: marker.Name(), Target: "/gitea-runner-probe", ReadOnly: true},
}},
})
if cerrdefs.IsInvalidArgument(err) {
return false, nil
@@ -100,24 +122,37 @@ func daemonSeesDir(ctx context.Context, dir string) (bool, error) {
if err != nil {
return false, err
}
_, err = cli.ContainerRemove(ctx, created.ID, client.ContainerRemoveOptions{Force: true})
return true, err
cleanupCtx, cancel := context.WithTimeout(context.WithoutCancel(ctx), dockerProxyProbeTimeout)
defer cancel()
if _, err := cli.ContainerRemove(cleanupCtx, created.ID, client.ContainerRemoveOptions{Force: true, RemoveVolumes: true}); err != nil {
return false, fmt.Errorf("removing the docker proxy probe container failed: %w", err)
}
return true, nil
}
// StartDockerProxy serves a job's docker socket in dir, labelling what the job creates through it.
func StartDockerProxy(daemonSocket, dir, job string) (*DockerProxy, error) {
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, err
}
digest := sha256.Sum256([]byte(job))
socket := filepath.Join(dir, hex.EncodeToString(digest[:8])+".sock")
_ = os.Remove(socket)
listener, err := net.Listen("unix", socket)
info, err := os.Stat(daemonSocket)
if err != nil {
return nil, err
}
if info, err := os.Stat(daemonSocket); err == nil {
_ = os.Chmod(socket, info.Mode().Perm())
if info.Mode()&os.ModeSocket == 0 {
return nil, errors.New("docker daemon path is not a Unix socket")
}
if err := os.MkdirAll(dir, 0o700); err != nil {
return nil, err
}
instance, err := os.MkdirTemp(dir, "p-")
if err != nil {
return nil, err
}
socket := filepath.Join(instance, "docker.sock")
listener, err := net.Listen("unix", socket)
if err != nil {
return nil, errors.Join(err, os.RemoveAll(instance))
}
if err := copyDockerSocketPermissions(socket, info); err != nil {
return nil, errors.Join(err, listener.Close(), os.RemoveAll(instance))
}
dial := func(ctx context.Context, _, _ string) (net.Conn, error) {
return (&net.Dialer{}).DialContext(ctx, "unix", daemonSocket)
@@ -130,61 +165,104 @@ func StartDockerProxy(daemonSocket, dir, job string) (*DockerProxy, error) {
},
Transport: transport,
}
server := &http.Server{ReadHeaderTimeout: 30 * time.Second, Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
switch {
case r.Method != http.MethodPost:
case createPath.MatchString(r.URL.Path):
streams, cancelStreams := context.WithCancel(context.Background())
creates, cancelCreates := context.WithCancel(context.Background())
var admission sync.Mutex
var handlers sync.WaitGroup
server := &http.Server{ReadHeaderTimeout: 30 * time.Second, ConnContext: func(ctx context.Context, conn net.Conn) context.Context {
return context.WithValue(ctx, dockerProxyConnKey{}, conn)
}, Handler: http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
admission.Lock()
if streams.Err() != nil {
admission.Unlock()
http.Error(w, "docker proxy is closing", http.StatusServiceUnavailable)
return
}
handlers.Add(1)
admission.Unlock()
defer handlers.Done()
creating := r.Method == http.MethodPost && createPath.MatchString(r.URL.Path)
parent, lifetime := r.Context(), streams
if creating {
parent, lifetime = context.WithoutCancel(parent), creates
}
ctx, cancel := context.WithCancel(parent)
defer cancel()
stop := context.AfterFunc(lifetime, func() {
cancel()
if !creating {
if conn, ok := parent.Value(dockerProxyConnKey{}).(net.Conn); ok {
_ = conn.Close()
}
}
})
defer stop()
r = r.WithContext(ctx)
if creating {
r.Body = http.MaxBytesReader(w, r.Body, maxCreateBody)
if err := addLabel(r, job); err != nil {
http.Error(w, err.Error(), http.StatusBadRequest)
status := http.StatusBadRequest
if _, ok := errors.AsType[*http.MaxBytesError](err); ok {
status = http.StatusRequestEntityTooLarge
}
http.Error(w, err.Error(), status)
return
}
case rawStreamPath.MatchString(r.URL.Path):
tunnel(w, r, dial)
} else if r.Method == http.MethodPost && rawStreamPath.MatchString(r.URL.Path) {
tunnel(w, r, dial, forward)
return
}
forward.ServeHTTP(w, r)
})}
go func() { _ = server.Serve(listener) }()
served := make(chan struct{})
go func() {
defer close(served)
_ = server.Serve(listener)
}()
return &DockerProxy{Socket: socket, close: func(ctx context.Context) error {
err := removeJobResources(ctx, job)
_ = server.Close()
ctx, cancel := context.WithTimeout(ctx, 30*time.Second)
defer cancel()
admission.Lock()
listenerErr := listener.Close()
cancelStreams()
admission.Unlock()
<-served
shutdownErr := server.Shutdown(ctx)
cancelCreates()
serverErr := server.Close()
handlers.Wait()
transport.CloseIdleConnections()
_ = os.Remove(socket)
return err
return errors.Join(ctx.Err(), listenerErr, shutdownErr, serverErr, os.RemoveAll(instance))
}}, nil
}
func addLabel(r *http.Request, job string) error {
body, err := io.ReadAll(io.LimitReader(r.Body, maxCreateBody+1))
body, err := io.ReadAll(r.Body)
if err != nil {
return err
}
if len(body) > maxCreateBody {
return errors.New("create request too large")
}
if len(bytes.TrimSpace(body)) == 0 {
r.Body = io.NopCloser(bytes.NewReader(body))
return nil
body = []byte("{}")
}
var fields map[string]json.RawMessage
var config struct{ Labels map[string]string }
if err := json.Unmarshal(body, &fields); err != nil {
return fmt.Errorf("invalid create request: %w", err)
}
key := "Labels"
for name := range fields {
if strings.EqualFold(name, key) {
key = name
break
}
if err := json.Unmarshal(body, &config); err != nil {
return fmt.Errorf("invalid create labels: %w", err)
}
labels := map[string]string{}
if raw := fields[key]; len(raw) > 0 && string(raw) != "null" {
if err := json.Unmarshal(raw, &labels); err != nil {
return fmt.Errorf("invalid create request: %w", err)
}
if fields == nil {
fields = make(map[string]json.RawMessage)
}
labels[jobLabel] = job
if fields[key], err = json.Marshal(labels); err != nil {
maps.DeleteFunc(fields, func(name string, _ json.RawMessage) bool {
return strings.EqualFold(name, "Labels")
})
if config.Labels == nil {
config.Labels = make(map[string]string)
}
config.Labels[jobLabel] = job
if fields["Labels"], err = json.Marshal(config.Labels); err != nil {
return err
}
if body, err = json.Marshal(fields); err != nil {
@@ -196,44 +274,92 @@ func addLabel(r *http.Request, job string) error {
return nil
}
type dockerProxyConnKey struct{}
type dockerProxyResponse struct {
response *http.Response
}
func (r dockerProxyResponse) RoundTrip(_ *http.Request) (*http.Response, error) {
return r.response, nil
}
// tunnel splices attach and exec streams, which the daemon hijacks with or without an HTTP upgrade
func tunnel(w http.ResponseWriter, r *http.Request, dial func(context.Context, string, string) (net.Conn, error)) {
hijacker, ok := w.(http.Hijacker)
if !ok {
http.Error(w, "connection cannot be hijacked", http.StatusInternalServerError)
return
}
func tunnel(w http.ResponseWriter, r *http.Request, dial func(context.Context, string, string) (net.Conn, error), forward *httputil.ReverseProxy) {
upstream, err := dial(r.Context(), "", "")
if err != nil {
http.Error(w, err.Error(), http.StatusBadGateway)
return
}
defer upstream.Close()
stop := context.AfterFunc(r.Context(), func() { _ = upstream.Close() })
defer stop()
if err := r.Write(upstream); err != nil {
http.Error(w, err.Error(), http.StatusBadGateway)
return
}
downstream, buffered, err := hijacker.Hijack()
reader := bufio.NewReader(upstream)
var response *http.Response
for {
response, err = http.ReadResponse(reader, r)
if err != nil {
http.Error(w, err.Error(), http.StatusBadGateway)
return
}
if response.StatusCode >= 200 || response.StatusCode == http.StatusSwitchingProtocols {
break
}
maps.Copy(w.Header(), response.Header)
w.WriteHeader(response.StatusCode)
clear(w.Header())
_ = response.Body.Close()
}
defer func() {
_ = upstream.Close()
_ = response.Body.Close()
}()
mediaType, _, _ := mime.ParseMediaType(response.Header.Get("Content-Type"))
if response.StatusCode != http.StatusSwitchingProtocols && (response.StatusCode != http.StatusOK || mediaType != "application/vnd.docker.raw-stream") {
ordinary := *forward
ordinary.Transport = dockerProxyResponse{response: response}
ordinary.ServeHTTP(w, r)
return
}
downstream, buffered, err := http.NewResponseController(w).Hijack()
if err != nil {
return
}
defer downstream.Close()
if _, err := io.CopyN(upstream, buffered, int64(buffered.Reader.Buffered())); err != nil {
if _, err := fmt.Fprintf(buffered, "%s %s\r\n", response.Proto, response.Status); err != nil {
return
}
done := make(chan struct{}, 2)
if err := response.Header.Write(buffered); err != nil {
return
}
if _, err := buffered.WriteString("\r\n"); err != nil {
return
}
if err := buffered.Flush(); err != nil {
return
}
done := make(chan struct{})
go func() {
_, _ = io.Copy(upstream, downstream)
done <- struct{}{}
}()
go func() {
_, _ = io.Copy(downstream, upstream)
done <- struct{}{}
defer close(done)
if _, err := io.Copy(upstream, io.MultiReader(io.LimitReader(buffered, int64(buffered.Reader.Buffered())), downstream)); err != nil { // Bypass net/http after the prefix so stdin EOF preserves output.
_ = upstream.Close()
} else if writer, ok := upstream.(interface{ CloseWrite() error }); ok {
_ = writer.CloseWrite()
} else {
_ = upstream.Close()
}
}()
_, _ = io.Copy(downstream, reader)
_ = downstream.Close()
_ = upstream.Close()
<-done
}
func removeJobResources(ctx context.Context, job string) error {
func RemoveDockerJobResources(ctx context.Context, job string) error {
cli, err := GetDockerClient(ctx)
if err != nil {
return err
@@ -246,27 +372,20 @@ func removeLabelled(ctx context.Context, cli client.APIClient, job string) error
logger := common.Logger(ctx)
filters := make(client.Filters).Add("label", jobLabel+"="+job)
containers, err := cli.ContainerList(ctx, client.ContainerListOptions{All: true, Filters: filters})
if err != nil {
return err
}
var errs []error
errs := []error{err}
for _, c := range containers.Items {
logger.Infof("removing container %s the job left behind", strings.TrimPrefix(strings.Join(c.Names, ","), "/"))
errs = append(errs, (&containerReference{cli: cli, id: c.ID}).remove()(ctx))
}
networks, err := cli.NetworkList(ctx, client.NetworkListOptions{Filters: filters})
if err != nil {
return errors.Join(append(errs, err)...)
}
errs = append(errs, err)
for _, n := range networks.Items {
if _, err := cli.NetworkRemove(ctx, n.ID, client.NetworkRemoveOptions{}); err != nil && !cerrdefs.IsNotFound(err) {
errs = append(errs, fmt.Errorf("failed to remove network %s: %w", n.Name, err))
}
}
volumes, err := cli.VolumeList(ctx, client.VolumeListOptions{Filters: filters})
if err != nil {
return errors.Join(append(errs, err)...)
}
errs = append(errs, err)
for _, v := range volumes.Items {
if _, err := cli.VolumeRemove(ctx, v.Name, client.VolumeRemoveOptions{}); err != nil && !cerrdefs.IsNotFound(err) {
errs = append(errs, fmt.Errorf("failed to remove volume %s: %w", v.Name, err))
+247 -48
View File
@@ -10,18 +10,23 @@ import (
"bufio"
"bytes"
"context"
"encoding/json"
"errors"
"fmt"
"io"
"net"
"net/http"
"os"
"path/filepath"
"runtime"
"strconv"
"strings"
"testing"
"time"
cerrdefs "github.com/containerd/errdefs"
"github.com/moby/moby/api/types/container"
"github.com/moby/moby/api/types/network"
"github.com/moby/moby/api/types/image"
"github.com/moby/moby/api/types/mount"
"github.com/moby/moby/api/types/volume"
mobyclient "github.com/moby/moby/client"
"github.com/stretchr/testify/assert"
@@ -56,7 +61,12 @@ func daemonSocketPath(t testing.TB, cli mobyclient.APIClient) string {
}
func TestDockerProxy(t *testing.T) {
if runtime.GOOS == "windows" {
t.Skip("Unix socket ownership is unavailable on Windows")
}
bodies := make(chan []byte, 8)
createStarted, releaseCreate := make(chan struct{}), make(chan struct{})
streamsDone := make(chan struct{}, 2)
daemonSocket := filepath.Join(shortTempDir(t), "d.sock")
listener, err := net.Listen("unix", daemonSocket)
require.NoError(t, err)
@@ -64,16 +74,33 @@ func TestDockerProxy(t *testing.T) {
switch {
case strings.HasSuffix(r.URL.Path, "/create"):
body, _ := io.ReadAll(r.Body)
if r.URL.RawQuery == "wait" {
close(createStarted)
<-releaseCreate
}
bodies <- body
w.WriteHeader(http.StatusCreated)
case strings.HasSuffix(r.URL.Path, "/start"):
case r.URL.Path == "/exec/detached/start" || r.URL.Path == "/exec/error/start":
if r.URL.Path == "/exec/error/start" {
w.Header().Set("Content-Type", "application/vnd.docker.raw-stream")
w.WriteHeader(http.StatusBadRequest)
}
_, _ = w.Write([]byte("ordinary"))
case rawStreamPath.MatchString(r.URL.Path) || r.URL.Path == "/session":
_, _ = io.Copy(io.Discard, r.Body)
conn, buffered, err := w.(http.Hijacker).Hijack()
if err != nil {
return
}
defer conn.Close()
_, _ = conn.Write([]byte("HTTP/1.1 200 OK\r\nContent-Type: application/vnd.docker.raw-stream\r\n\r\n"))
_, _ = io.Copy(conn, buffered)
if r.Header.Get("Upgrade") != "" {
_, _ = conn.Write([]byte("HTTP/1.1 101 Switching Protocols\r\nConnection: Upgrade\r\nUpgrade: tcp\r\n\r\nready\n"))
} else {
_, _ = conn.Write([]byte("HTTP/1.1 200 OK\r\nContent-Type: application/vnd.docker.raw-stream\r\n\r\nready\n"))
}
input, _ := io.ReadAll(buffered)
_, _ = conn.Write(append([]byte("final:"), input...))
streamsDone <- struct{}{}
default:
w.Header().Set("Api-Version", "1.47")
_, _ = w.Write([]byte("OK " + r.Method + " " + r.URL.Path))
@@ -83,34 +110,69 @@ func TestDockerProxy(t *testing.T) {
t.Cleanup(func() { _ = daemon.Close() })
proxy, err := StartDockerProxy(daemonSocket, shortTempDir(t), "job-1")
require.NoError(t, err)
client := &http.Client{Transport: &http.Transport{DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
t.Cleanup(func() { _ = proxy.Close(context.Background()) })
client := &http.Client{Timeout: 5 * time.Second, Transport: &http.Transport{DialContext: func(ctx context.Context, _, _ string) (net.Conn, error) {
return (&net.Dialer{}).DialContext(ctx, "unix", proxy.Socket)
}}}
t.Cleanup(client.CloseIdleConnections)
dialProxy := func(t *testing.T) *net.UnixConn {
t.Helper()
conn, err := net.DialUnix("unix", nil, &net.UnixAddr{Name: proxy.Socket, Net: "unix"})
require.NoError(t, err)
t.Cleanup(func() { _ = conn.Close() })
require.NoError(t, conn.SetDeadline(time.Now().Add(5*time.Second)))
return conn
}
startStream := func(t *testing.T, path, upgrade string) (*net.UnixConn, *bufio.Reader) {
t.Helper()
conn := dialProxy(t)
_, err := fmt.Fprintf(conn, "POST %s HTTP/1.1\r\nHost: docker\r\n%sContent-Length: 2\r\n\r\n{}input", path, upgrade)
require.NoError(t, err)
reader := bufio.NewReader(conn)
resp, err := http.ReadResponse(reader, nil)
require.NoError(t, err)
t.Cleanup(func() {
_ = conn.Close()
_ = resp.Body.Close()
})
prefix, err := reader.ReadString('\n')
require.NoError(t, err)
require.Equal(t, "ready\n", prefix)
return conn, reader
}
t.Run("labels creates", func(t *testing.T) {
for path, body := range map[string]string{
"/v1.47/containers/create": `{"Image":"alpine","Labels":{"own":"1"}}`,
"/networks/create": `{"Name":"n"}`,
"/volumes/create": `{"Name":"v","labels":null}`,
for _, testCase := range []struct{ path, body, want string }{
{
"/v1.47/containers/create", `{"Image":"alpine","Unknown":{"enabled":true},"Labels":{"own":"1","com.gitea.runner.job":"other"}}`,
`{"Image":"alpine","Unknown":{"enabled":true},"Labels":{"own":"1","com.gitea.runner.job":"job-1"}}`,
},
{
"/networks/create", `{"Name":"n"}`,
`{"Name":"n","Labels":{"com.gitea.runner.job":"job-1"}}`,
},
{
"/volumes/create", `{"Name":"v","labels":null}`,
`{"Name":"v","Labels":{"com.gitea.runner.job":"job-1"}}`,
},
{
"/volumes/create", "",
`{"Labels":{"com.gitea.runner.job":"job-1"}}`,
},
{
"/volumes/create", "null",
`{"Labels":{"com.gitea.runner.job":"job-1"}}`,
},
{
"/volumes/create", `{"Labels":{"discard":"1"},"labels":null,"LABELS":{"own":"1"},"LABELS":{"extra":"2"}}`,
`{"Labels":{"own":"1","extra":"2","com.gitea.runner.job":"job-1"}}`,
},
} {
resp, err := client.Post("http://docker"+path, "application/json", strings.NewReader(body))
resp, err := client.Post("http://docker"+testCase.path, "application/json", strings.NewReader(testCase.body))
require.NoError(t, err)
resp.Body.Close()
assert.Equal(t, http.StatusCreated, resp.StatusCode)
var got map[string]json.RawMessage
require.NoError(t, json.Unmarshal(<-bodies, &got))
key := "Labels"
if _, ok := got["labels"]; ok {
key = "labels"
}
var labels map[string]string
require.NoError(t, json.Unmarshal(got[key], &labels))
assert.Equal(t, "job-1", labels[jobLabel], path)
if strings.Contains(body, "own") {
assert.Equal(t, "1", labels["own"])
assert.JSONEq(t, `"alpine"`, string(got["Image"]))
}
require.Equal(t, http.StatusCreated, resp.StatusCode)
assert.JSONEq(t, testCase.want, string(<-bodies), testCase.body)
}
})
@@ -123,40 +185,99 @@ func TestDockerProxy(t *testing.T) {
assert.Equal(t, "OK GET /v1.47/_ping", string(body))
})
t.Run("tunnels raw streams", func(t *testing.T) {
conn, err := net.Dial("unix", proxy.Socket)
require.NoError(t, err)
defer conn.Close()
_, err = conn.Write([]byte("POST /v1.47/exec/abc/start HTTP/1.1\r\nHost: docker\r\nContent-Length: 2\r\n\r\n{}ping\n"))
require.NoError(t, err)
t.Run("tunnels raw streams through stdin EOF", func(t *testing.T) {
for _, stream := range []struct{ path, upgrade string }{
{path: "/v1.47/exec/abc/start"},
{path: "/containers/abc/attach", upgrade: "Connection: Upgrade\r\nUpgrade: tcp\r\n"},
} {
conn, reader := startStream(t, stream.path, stream.upgrade)
require.NoError(t, conn.CloseWrite())
output, err := io.ReadAll(reader)
require.NoError(t, err)
assert.Equal(t, "final:input", string(output))
<-streamsDone
}
})
t.Run("detached and error responses retain keepalive label injection", func(t *testing.T) {
conn := dialProxy(t)
reader := bufio.NewReader(conn)
resp, err := http.ReadResponse(reader, nil)
for endpoint, status := range map[string]int{"detached": http.StatusOK, "error": http.StatusBadRequest} {
_, err := fmt.Fprintf(conn, "POST /exec/%s/start HTTP/1.1\r\nHost: docker\r\nContent-Length: 2\r\n\r\n{}", endpoint)
require.NoError(t, err)
resp, err := http.ReadResponse(reader, nil)
require.NoError(t, err)
body, err := io.ReadAll(resp.Body)
require.NoError(t, err)
require.NoError(t, resp.Body.Close())
assert.Equal(t, status, resp.StatusCode)
assert.Equal(t, "ordinary", string(body))
_, err = io.WriteString(conn, "POST /volumes/create HTTP/1.1\r\nHost: docker\r\nContent-Length: 0\r\n\r\n")
require.NoError(t, err)
resp, err = http.ReadResponse(reader, nil)
require.NoError(t, err)
require.NoError(t, resp.Body.Close())
require.Equal(t, http.StatusCreated, resp.StatusCode)
assert.JSONEq(t, `{"Labels":{"com.gitea.runner.job":"job-1"}}`, string(<-bodies))
}
})
t.Run("close joins streams and admitted create", func(t *testing.T) {
t.Cleanup(func() { close(releaseCreate) })
_, raw := startStream(t, "/exec/live/start", "")
_, upgraded := startStream(t, "/session", "Connection: Upgrade\r\nUpgrade: tcp\r\n")
creator := dialProxy(t)
_, err := io.WriteString(creator, "POST /volumes/create?wait HTTP/1.1\r\nHost: docker\r\nContent-Length: 2\r\n\r\n{}")
require.NoError(t, err)
defer resp.Body.Close()
assert.Equal(t, "application/vnd.docker.raw-stream", resp.Header.Get("Content-Type"))
echoed, err := reader.ReadString('\n')
<-createStarted
require.NoError(t, creator.CloseWrite())
closed := make(chan error, 1)
ctx, cancel := context.WithTimeout(t.Context(), 5*time.Second)
defer cancel()
go func() { closed <- proxy.Close(ctx) }()
for _, reader := range []*bufio.Reader{raw, upgraded} {
_, err := reader.ReadByte()
require.ErrorIs(t, err, io.EOF)
<-streamsDone
}
select {
case err := <-closed:
t.Fatalf("Close returned before admitted create settled: %v", err)
default:
}
releaseCreate <- struct{}{}
resp, err := http.ReadResponse(bufio.NewReader(creator), nil)
require.NoError(t, err)
assert.Equal(t, "{}ping\n", echoed)
require.NoError(t, resp.Body.Close())
require.Equal(t, http.StatusCreated, resp.StatusCode)
assert.JSONEq(t, `{"Labels":{"com.gitea.runner.job":"job-1"}}`, string(<-bodies))
require.NoError(t, <-closed)
})
}
func TestRemoveLabelledRemovesContainersNetworksAndVolumes(t *testing.T) {
containerFailure := errors.New("container removal failed")
listFailure := errors.New("listing failed")
volumeFailure := errors.New("volume removal failed")
ctx := context.Background()
filters := make(mobyclient.Filters).Add("label", jobLabel+"=job-1")
cli := &mockDockerClient{}
cli.On("ContainerList", ctx, mobyclient.ContainerListOptions{All: true, Filters: filters}).
Return(mobyclient.ContainerListResult{Items: []container.Summary{{ID: "c1", Names: []string{"/app"}}}}, nil)
cli.On("ContainerKill", ctx, "c1", mock.Anything).Return(mobyclient.ContainerKillResult{}, nil)
Return(mobyclient.ContainerListResult{Items: []container.Summary{{ID: "c1", Names: []string{"/app"}}}}, nil).Once()
cli.On("ContainerKill", ctx, "c1", mock.Anything).Return(mobyclient.ContainerKillResult{}, nil).Once()
cli.On("ContainerRemove", ctx, "c1", mobyclient.ContainerRemoveOptions{RemoveVolumes: true, Force: true}).
Return(mobyclient.ContainerRemoveResult{}, nil)
cli.On("NetworkList", ctx, mobyclient.NetworkListOptions{Filters: filters}).
Return(mobyclient.NetworkListResult{Items: []network.Summary{{ID: "n1", Name: "app_default"}}}, nil)
cli.On("NetworkRemove", ctx, "n1", mock.Anything).Return(mobyclient.NetworkRemoveResult{}, nil)
Return(mobyclient.ContainerRemoveResult{}, containerFailure).Once()
cli.On("NetworkList", ctx, mobyclient.NetworkListOptions{Filters: filters}).Return(mobyclient.NetworkListResult{}, listFailure).Once()
cli.On("VolumeList", ctx, mobyclient.VolumeListOptions{Filters: filters}).
Return(mobyclient.VolumeListResult{Items: []volume.Volume{{Name: "app_data"}}}, nil)
cli.On("VolumeRemove", ctx, "app_data", mobyclient.VolumeRemoveOptions{}).Return(mobyclient.VolumeRemoveResult{}, nil)
Return(mobyclient.VolumeListResult{Items: []volume.Volume{{Name: "app_data"}}}, nil).Once()
cli.On("VolumeRemove", ctx, "app_data", mobyclient.VolumeRemoveOptions{}).Return(mobyclient.VolumeRemoveResult{}, volumeFailure).Once()
require.NoError(t, removeLabelled(ctx, cli, "job-1"))
err := removeLabelled(ctx, cli, "job-1")
require.ErrorIs(t, err, containerFailure)
require.ErrorIs(t, err, listFailure)
require.ErrorIs(t, err, volumeFailure)
require.ErrorContains(t, err, "failed to remove container c1")
require.ErrorContains(t, err, "failed to remove volume app_data")
cli.AssertExpectations(t)
}
@@ -168,13 +289,14 @@ func TestDockerProxyWithDaemon(t *testing.T) {
require.NoError(t, err)
defer direct.Close()
dir := shortTempDir(t)
seen, err := daemonSeesDir(ctx, dir)
seen, err := daemonSeesDir(ctx, direct, dir)
require.NoError(t, err)
t.Logf("daemon sees the runner's filesystem: %v", seen)
job := "proxy-test-" + t.Name()
job := "proxy-test-" + strconv.FormatInt(time.Now().UnixNano(), 36)
proxy, err := StartDockerProxy(daemonSocketPath(t, direct), dir, job)
require.NoError(t, err)
t.Cleanup(func() { _ = proxy.Close(context.Background()) })
viaProxy, err := mobyclient.New(mobyclient.WithHost("unix://" + proxy.Socket))
require.NoError(t, err)
defer viaProxy.Close()
@@ -212,6 +334,7 @@ func TestDockerProxyWithDaemon(t *testing.T) {
attached.Close()
require.NoError(t, proxy.Close(ctx))
require.NoError(t, RemoveDockerJobResources(ctx, job))
_, err = direct.ContainerInspect(ctx, created.ID, mobyclient.ContainerInspectOptions{})
assert.True(t, cerrdefs.IsNotFound(err))
_, err = direct.NetworkInspect(ctx, net.ID, mobyclient.NetworkInspectOptions{})
@@ -268,3 +391,79 @@ func BenchmarkDockerProxy(b *testing.B) {
})
}
}
type probeClient struct {
mobyclient.APIClient
create func(mobyclient.ContainerCreateOptions) (mobyclient.ContainerCreateResult, error)
remove func(context.Context, string, mobyclient.ContainerRemoveOptions) (mobyclient.ContainerRemoveResult, error)
}
func (c *probeClient) ImageList(context.Context, mobyclient.ImageListOptions) (mobyclient.ImageListResult, error) {
return mobyclient.ImageListResult{Items: []image.Summary{{ID: "probe-image"}}}, nil
}
func (c *probeClient) ContainerCreate(_ context.Context, opts mobyclient.ContainerCreateOptions) (mobyclient.ContainerCreateResult, error) {
return c.create(opts)
}
func (c *probeClient) ContainerRemove(ctx context.Context, id string, opts mobyclient.ContainerRemoveOptions) (mobyclient.ContainerRemoveResult, error) {
return c.remove(ctx, id, opts)
}
func TestDaemonSeesDir(t *testing.T) {
dir := t.TempDir()
markers := make(map[string]bool)
for _, testCase := range []struct {
name string
private bool
removeErr error
}{
{name: "cleanup after cancellation"},
{name: "private filesystem with stale marker", private: true},
{name: "cleanup error after cancellation", removeErr: errors.New("cleanup failed")},
} {
t.Run(testCase.name, func(t *testing.T) {
ctx, cancel := context.WithCancel(t.Context())
defer cancel()
daemonDir := dir
if testCase.private {
daemonDir = t.TempDir()
require.NoError(t, os.WriteFile(filepath.Join(daemonDir, "probe"), []byte("stale"), 0o600))
}
removed := false
cli := &probeClient{
create: func(opts mobyclient.ContainerCreateOptions) (mobyclient.ContainerCreateResult, error) {
require.Len(t, opts.HostConfig.Mounts, 1)
marker := opts.HostConfig.Mounts[0].Source
assert.Equal(t, mount.Mount{Type: mount.TypeBind, Source: marker, Target: "/gitea-runner-probe", ReadOnly: true}, opts.HostConfig.Mounts[0])
assert.False(t, markers[marker])
markers[marker] = true
info, err := os.Stat(marker)
require.NoError(t, err)
assert.True(t, info.Mode().IsRegular())
cancel()
if _, err := os.Stat(filepath.Join(daemonDir, filepath.Base(marker))); errors.Is(err, os.ErrNotExist) {
return mobyclient.ContainerCreateResult{}, cerrdefs.ErrInvalidArgument
}
return mobyclient.ContainerCreateResult{ID: "probe"}, nil
},
remove: func(ctx context.Context, id string, opts mobyclient.ContainerRemoveOptions) (mobyclient.ContainerRemoveResult, error) {
removed = true
assert.Equal(t, "probe", id)
assert.Equal(t, mobyclient.ContainerRemoveOptions{Force: true, RemoveVolumes: true}, opts)
require.NoError(t, ctx.Err())
_, bounded := ctx.Deadline()
assert.True(t, bounded)
return mobyclient.ContainerRemoveResult{}, testCase.removeErr
},
}
seen, err := daemonSeesDir(ctx, cli, dir)
require.ErrorIs(t, err, testCase.removeErr)
assert.Equal(t, !testCase.private && testCase.removeErr == nil, seen)
assert.Equal(t, !testCase.private, removed)
entries, err := os.ReadDir(dir)
require.NoError(t, err)
assert.Empty(t, entries)
})
}
}
+27
View File
@@ -0,0 +1,27 @@
// Copyright 2026 The Gitea Authors. All rights reserved.
// SPDX-License-Identifier: MIT
//go:build !WITHOUT_DOCKER && (linux || darwin || netbsd)
package container
import (
"errors"
"os"
"path/filepath"
"syscall"
)
func copyDockerSocketPermissions(socket string, info os.FileInfo) error {
stat, ok := info.Sys().(*syscall.Stat_t)
if !ok {
return errors.New("docker socket ownership is unavailable")
}
if err := os.Chown(socket, int(stat.Uid), int(stat.Gid)); err != nil {
return err
}
if err := os.Chown(filepath.Dir(socket), int(stat.Uid), -1); err != nil {
return err
}
return os.Chmod(socket, info.Mode().Perm())
}
+15
View File
@@ -0,0 +1,15 @@
// Copyright 2026 The Gitea Authors. All rights reserved.
// SPDX-License-Identifier: MIT
//go:build !WITHOUT_DOCKER
package container
import (
"errors"
"os"
)
func copyDockerSocketPermissions(_ string, _ os.FileInfo) error {
return errors.New("docker socket ownership cannot be preserved on Windows")
}
+13 -21
View File
@@ -488,14 +488,11 @@ func (cr *containerReference) remove() common.Executor {
RemoveVolumes: true,
Force: true,
})
switch {
case cerrdefs.IsConflict(err):
// the daemon's own AutoRemove teardown is running, and it releases the volume
// references and the network endpoint only once it finishes
cr.waitForRemoval(ctx, idOrName)
case err != nil && !cerrdefs.IsNotFound(err):
logger.Error(fmt.Errorf("failed to remove container %s: %w", idOrName, err))
return nil // keep the id, the container is still there for a later Remove()
if cerrdefs.IsConflict(err) {
err = cr.waitForRemoval(ctx, idOrName)
}
if err != nil && !cerrdefs.IsNotFound(err) {
return fmt.Errorf("failed to remove container %s: %w", idOrName, err)
}
logger.Debugf("Removed container: %v", idOrName)
@@ -504,7 +501,7 @@ func (cr *containerReference) remove() common.Executor {
}
}
func (cr *containerReference) waitForRemoval(ctx context.Context, idOrName string) {
func (cr *containerReference) waitForRemoval(ctx context.Context, idOrName string) error {
// per container, against the one minute the post-job executor allows for the whole
// cleanup, so a job with several services can spend most of that budget here
ctx, cancel := context.WithTimeout(ctx, 15*time.Second)
@@ -514,18 +511,13 @@ func (cr *containerReference) waitForRemoval(ctx context.Context, idOrName strin
Condition: container.WaitConditionRemoved,
})
select {
case <-waitResult.Result:
case <-waitResult.Error:
case <-ctx.Done():
// the client delivers the result over an unbuffered channel, so leave a receiver
// behind or its goroutine parks on the send for the lifetime of the process
go func() {
select {
case <-waitResult.Result:
case <-waitResult.Error:
}
}()
common.Logger(ctx).Warnf("Timed out waiting for the daemon to remove container %s, its volumes and network may be left behind", idOrName)
case result := <-waitResult.Result:
if result.Error != nil {
return errors.New(result.Error.Message)
}
return nil
case err := <-waitResult.Error:
return err
}
}
+19 -12
View File
@@ -345,18 +345,22 @@ func TestDockerWaitFailure(t *testing.T) {
// A remove that raced the daemon's AutoRemove teardown is not a failure and must not
// be logged as one.
func TestRemoveIgnoresAutoRemoveRace(t *testing.T) {
removeFailure := errors.New("driver failed to remove root filesystem")
removeOpts := mobyclient.ContainerRemoveOptions{RemoveVolumes: true, Force: true}
killOpts := mobyclient.ContainerKillOptions{Signal: "SIGKILL"}
for _, tc := range []struct {
name string
err error
wantWait bool
wantFailure bool
name string
err error
wantWait bool
waitErr error
wantErr error
}{
{name: "removal in progress", err: cerrdefs.ErrConflict.WithMessage("removal of container abc is already in progress"), wantWait: true},
{name: "wait canceled", err: cerrdefs.ErrConflict, wantWait: true, waitErr: context.Canceled, wantErr: context.Canceled},
{name: "removed during wait", err: cerrdefs.ErrConflict, wantWait: true, waitErr: cerrdefs.ErrNotFound},
{name: "already removed", err: cerrdefs.ErrNotFound.WithMessage("No such container: abc")},
{name: "removed cleanly", err: nil},
{name: "real failure", err: errors.New("driver failed to remove root filesystem"), wantFailure: true},
{name: "real failure", err: removeFailure, wantErr: removeFailure},
} {
t.Run(tc.name, func(t *testing.T) {
logger, hook := test.NewNullLogger()
@@ -366,21 +370,24 @@ func TestRemoveIgnoresAutoRemoveRace(t *testing.T) {
client.On("ContainerRemove", ctx, "abc", removeOpts).Return(mobyclient.ContainerRemoveResult{}, tc.err)
if tc.wantWait {
removed := make(chan container.WaitResponse, 1)
removed <- container.WaitResponse{}
waitErrors := make(chan error, 1)
if tc.waitErr == nil {
removed <- container.WaitResponse{}
} else {
waitErrors <- tc.waitErr
}
client.On("ContainerWait", mock.Anything, "abc", mobyclient.ContainerWaitOptions{Condition: container.WaitConditionRemoved}).
Return(mobyclient.ContainerWaitResult{Result: removed})
Return(mobyclient.ContainerWaitResult{Result: removed, Error: waitErrors})
}
cr := &containerReference{id: "abc", cli: client}
require.NoError(t, cr.remove()(ctx))
// a failure keeps the id, so a later Remove() can retry it
if tc.wantFailure {
require.ErrorIs(t, cr.remove()(ctx), tc.wantErr)
if tc.wantErr != nil {
assert.Equal(t, "abc", cr.id)
assert.Len(t, hook.AllEntries(), 1)
} else {
assert.Empty(t, cr.id)
assert.Empty(t, hook.AllEntries())
}
assert.Empty(t, hook.AllEntries())
client.AssertExpectations(t)
})
}
+6 -2
View File
@@ -29,8 +29,12 @@ func RemoveImage(ctx context.Context, imageName string, force, pruneChildren boo
return false, errors.New("Unsupported Operation")
}
func DockerProxyDir(ctx context.Context) string {
return ""
func RemoveDockerJobResources(_ context.Context, _ string) error {
return nil
}
func NewDockerProxy(_ context.Context, _ string) *DockerProxy {
return nil
}
func StartDockerProxy(daemonSocket, dir, job string) (*DockerProxy, error) {