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
9 changes: 9 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,15 @@ backend:

This setting does not control how you get the long-running `oz-agent-worker` image (see [Image releases and pinning](#image-releases-and-pinning)). It controls only the task and sidecar images the worker uses to run tasks.

`backend.docker.sidecar_image` overrides the warp-agent sidecar image reference sent by the server (e.g. `docker.io/warpdotdev/warp-agent:latest`); set this when the worker host cannot pull directly from Docker Hub and must use an internal registry mirror or pull-through cache instead. This only affects the warp-agent sidecar (mounted at `/agent`), not any additional sidecars. When using this override, you are responsible for keeping your mirror in sync with `docker.io/warpdotdev/warp-agent` — the server normally sends the correct version-matched image per task, so a stale mirror may cause version incompatibility.

```yaml
worker_id: "my-worker"
backend:
docker:
sidecar_image: "my-registry.io/warpdotdev/warp-agent:latest"
```

### Image releases and pinning

Production self-hosted workers should pin an immutable image version instead of relying on `latest`.
Expand Down
1 change: 1 addition & 0 deletions internal/config/config.go
Original file line number Diff line number Diff line change
Expand Up @@ -56,6 +56,7 @@ type DockerConfig struct {
// Kubernetes backend, an omitted value defaults to Always here, preserving the Docker
// backend's original unconditional-pull behavior for existing installations.
ImagePullPolicy string `yaml:"image_pull_policy" validate:"omitempty,oneof=Always Never IfNotPresent"`
SidecarImage string `yaml:"sidecar_image" validate:"omitempty,no_whitespace"`
Environment []EnvEntry `yaml:"environment" validate:"dive"`
}

Expand Down
50 changes: 50 additions & 0 deletions internal/config/config_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -646,6 +646,56 @@ backend:
})
}

func TestLoadDockerSidecarImage(t *testing.T) {
t.Run("parses sidecar_image when set", func(t *testing.T) {
path := writeTestConfig(t, `
worker_id: "docker-worker"
backend:
docker:
sidecar_image: "my-registry.io/warpdotdev/warp-agent:latest"
`)
cfg, err := Load(path)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if cfg.Backend.Docker == nil {
t.Fatal("expected docker backend to be set")
}
if cfg.Backend.Docker.SidecarImage != "my-registry.io/warpdotdev/warp-agent:latest" {
t.Errorf("sidecar_image = %q, want %q", cfg.Backend.Docker.SidecarImage, "my-registry.io/warpdotdev/warp-agent:latest")
}
})

t.Run("sidecar_image is empty when not set", func(t *testing.T) {
path := writeTestConfig(t, `
worker_id: "docker-worker"
backend:
docker:
volumes: []
`)
cfg, err := Load(path)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if cfg.Backend.Docker.SidecarImage != "" {
t.Errorf("expected sidecar_image to be empty, got %q", cfg.Backend.Docker.SidecarImage)
}
})

t.Run("rejects sidecar_image with whitespace", func(t *testing.T) {
path := writeTestConfig(t, `
worker_id: "docker-worker"
backend:
docker:
sidecar_image: "my image:latest"
`)
_, err := Load(path)
if err == nil {
t.Fatal("expected error for sidecar_image with whitespace")
}
})
}

func TestLoadKubernetesSidecarImage(t *testing.T) {
t.Run("parses sidecar_image when set", func(t *testing.T) {
path := writeTestConfig(t, `
Expand Down
1 change: 1 addition & 0 deletions internal/worker/docker.go
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,7 @@ type DockerBackendConfig struct {
// empty value defaults to PullPolicyAlways, preserving the backend's original
// unconditional-pull behavior for installations that don't set it.
ImagePullPolicy string
SidecarImage string
}

func (b *DockerBackend) containerWasOOMKilled(ctx context.Context, dockerClient *client.Client, containerID string) bool {
Expand Down
19 changes: 16 additions & 3 deletions internal/worker/worker.go
Original file line number Diff line number Diff line change
Expand Up @@ -533,9 +533,9 @@ func (w *Worker) prepareTaskParams(assignment *types.TaskAssignmentMessage) *Tas
var sidecars []types.SidecarMount
if assignment.SidecarImage != "" {
sidecarImage := assignment.SidecarImage
if w.config.Kubernetes != nil && w.config.Kubernetes.SidecarImage != "" {
log.Infof(w.ctx, "Overriding server sidecar image %s with configured sidecar image %s", assignment.SidecarImage, w.config.Kubernetes.SidecarImage)
sidecarImage = w.config.Kubernetes.SidecarImage
if override := w.configuredWarpAgentSidecarImage(); override != "" {
log.Infof(w.ctx, "Overriding server sidecar image %s with configured sidecar image %s", assignment.SidecarImage, override)
sidecarImage = override
}
sidecars = append(sidecars, types.SidecarMount{
Image: sidecarImage,
Expand Down Expand Up @@ -600,6 +600,19 @@ func (w *Worker) prepareTaskParams(assignment *types.TaskAssignmentMessage) *Tas
}
}

// configuredWarpAgentSidecarImage returns the operator-configured warp-agent sidecar
// image, or empty if neither backend set one. Only one backend config is populated
// at runtime.
func (w *Worker) configuredWarpAgentSidecarImage() string {
if w.config.Kubernetes != nil && w.config.Kubernetes.SidecarImage != "" {
return w.config.Kubernetes.SidecarImage
}
if w.config.Docker != nil && w.config.Docker.SidecarImage != "" {
return w.config.Docker.SidecarImage
}
return ""
}

// defaultImageForTask returns the Docker image to use for a task, applying the
// precedence: server-provided > worker config default_image > hardcoded fallback.
func (w *Worker) defaultImageForTask(assignmentImage string, task *types.Task) string {
Expand Down
111 changes: 74 additions & 37 deletions internal/worker/worker_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -622,63 +622,100 @@ func TestDefaultImageForTask(t *testing.T) {
}

func TestPrepareTaskParamsSidecarImageOverride(t *testing.T) {
newWorker := func(sidecarImage string) *Worker {
ctx := context.Background()
var k8sConfig *KubernetesBackendConfig
if sidecarImage != "" {
k8sConfig = &KubernetesBackendConfig{SidecarImage: sidecarImage}
} else {
k8sConfig = &KubernetesBackendConfig{}
const (
serverImage = "docker.io/warpdotdev/warp-agent:latest"
overrideImage = "my-registry.io/warpdotdev/warp-agent:latest"
)

newK8sWorker := func(sidecarImage string) *Worker {
return &Worker{
ctx: context.Background(),
config: Config{
Kubernetes: &KubernetesBackendConfig{SidecarImage: sidecarImage},
},
}
}
newDockerWorker := func(sidecarImage string) *Worker {
return &Worker{
ctx: ctx,
ctx: context.Background(),
config: Config{
Kubernetes: k8sConfig,
Docker: &DockerBackendConfig{SidecarImage: sidecarImage},
},
}
}
assignment := func(sidecarImage string, extra ...types.SidecarMount) *types.TaskAssignmentMessage {
return &types.TaskAssignmentMessage{
TaskID: "task-1",
Task: &types.Task{ID: "task-1"},
SidecarImage: sidecarImage,
AdditionalSidecars: extra,
}
}

t.Run("config sidecar_image overrides server-provided image", func(t *testing.T) {
w := newWorker("my-registry.io/warpdotdev/warp-agent:latest")
params := w.prepareTaskParams(&types.TaskAssignmentMessage{
TaskID: "task-1",
Task: &types.Task{ID: "task-1"},
SidecarImage: "docker.io/warpdotdev/warp-agent:latest",
})
assertAgentSidecar := func(t *testing.T, params *TaskParams, wantImage string) {
t.Helper()
if len(params.Sidecars) == 0 {
t.Fatal("expected at least one sidecar")
}
if params.Sidecars[0].Image != "my-registry.io/warpdotdev/warp-agent:latest" {
t.Errorf("sidecar image = %q, want %q", params.Sidecars[0].Image, "my-registry.io/warpdotdev/warp-agent:latest")
if params.Sidecars[0].Image != wantImage {
t.Errorf("sidecar image = %q, want %q", params.Sidecars[0].Image, wantImage)
}
if params.Sidecars[0].MountPath != "/agent" {
t.Errorf("sidecar mount path = %q, want /agent", params.Sidecars[0].MountPath)
}
}

t.Run("kubernetes config sidecar_image overrides server-provided image", func(t *testing.T) {
params := newK8sWorker(overrideImage).prepareTaskParams(assignment(serverImage))
assertAgentSidecar(t, params, overrideImage)
})

t.Run("server-provided image used when config sidecar_image empty", func(t *testing.T) {
w := newWorker("")
params := w.prepareTaskParams(&types.TaskAssignmentMessage{
TaskID: "task-1",
Task: &types.Task{ID: "task-1"},
SidecarImage: "docker.io/warpdotdev/warp-agent:latest",
})
if len(params.Sidecars) == 0 {
t.Fatal("expected at least one sidecar")
}
if params.Sidecars[0].Image != "docker.io/warpdotdev/warp-agent:latest" {
t.Errorf("sidecar image = %q, want %q", params.Sidecars[0].Image, "docker.io/warpdotdev/warp-agent:latest")
t.Run("docker config sidecar_image overrides server-provided image", func(t *testing.T) {
params := newDockerWorker(overrideImage).prepareTaskParams(assignment(serverImage))
assertAgentSidecar(t, params, overrideImage)
})

t.Run("server-provided image used when kubernetes config sidecar_image empty", func(t *testing.T) {
params := newK8sWorker("").prepareTaskParams(assignment(serverImage))
assertAgentSidecar(t, params, serverImage)
})

t.Run("server-provided image used when docker config sidecar_image empty", func(t *testing.T) {
params := newDockerWorker("").prepareTaskParams(assignment(serverImage))
assertAgentSidecar(t, params, serverImage)
})

t.Run("server-provided image used when docker config is nil", func(t *testing.T) {
w := &Worker{ctx: context.Background(), config: Config{}}
params := w.prepareTaskParams(assignment(serverImage))
assertAgentSidecar(t, params, serverImage)
})

t.Run("no sidecar when server provides empty sidecar image with kubernetes override", func(t *testing.T) {
params := newK8sWorker(overrideImage).prepareTaskParams(assignment(""))
if len(params.Sidecars) != 0 {
t.Errorf("expected no sidecars when server sidecar image is empty, got %d", len(params.Sidecars))
}
})

t.Run("no sidecar when server provides empty sidecar image", func(t *testing.T) {
w := newWorker("my-registry.io/warpdotdev/warp-agent:latest")
params := w.prepareTaskParams(&types.TaskAssignmentMessage{
TaskID: "task-1",
Task: &types.Task{ID: "task-1"},
SidecarImage: "",
})
t.Run("no sidecar when server provides empty sidecar image with docker override", func(t *testing.T) {
params := newDockerWorker(overrideImage).prepareTaskParams(assignment(""))
if len(params.Sidecars) != 0 {
t.Errorf("expected no sidecars when server sidecar image is empty, got %d", len(params.Sidecars))
}
})

t.Run("docker override does not rewrite additional sidecars", func(t *testing.T) {
extra := types.SidecarMount{Image: "docker.io/warpdotdev/extra:latest", MountPath: "/mnt/extra"}
params := newDockerWorker(overrideImage).prepareTaskParams(assignment(serverImage, extra))
if len(params.Sidecars) != 2 {
t.Fatalf("sidecar count = %d, want 2: %+v", len(params.Sidecars), params.Sidecars)
}
assertAgentSidecar(t, params, overrideImage)
if params.Sidecars[1] != extra {
t.Errorf("additional sidecar = %+v, want %+v", params.Sidecars[1], extra)
}
})
}

func TestPrepareTaskParamsCodingCLISidecarOverride(t *testing.T) {
Expand Down
4 changes: 3 additions & 1 deletion main.go
Original file line number Diff line number Diff line change
Expand Up @@ -362,10 +362,11 @@ func mergeConfig(fileConfig *config.FileConfig) (worker.Config, error) {

// Merge volumes: config file + CLI (concatenated).
var volumes []string
var imagePullPolicy string
var imagePullPolicy, sidecarImage string
if fileConfig != nil && fileConfig.Backend.Docker != nil {
volumes = append(volumes, fileConfig.Backend.Docker.Volumes...)
imagePullPolicy = fileConfig.Backend.Docker.ImagePullPolicy
sidecarImage = fileConfig.Backend.Docker.SidecarImage
}
volumes = append(volumes, CLI.Volumes...)

Expand All @@ -374,6 +375,7 @@ func mergeConfig(fileConfig *config.FileConfig) (worker.Config, error) {
Volumes: volumes,
Env: mergedEnv,
ImagePullPolicy: imagePullPolicy,
SidecarImage: sidecarImage,
}
}

Expand Down
52 changes: 52 additions & 0 deletions main_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -298,6 +298,58 @@ func TestMergeConfigDockerImagePullPolicyFromFile(t *testing.T) {
}
}

func TestMergeConfigDockerSidecarImageFromFile(t *testing.T) {
resetCLIForTest()
t.Cleanup(resetCLIForTest)

fileConfig := &config.FileConfig{
WorkerID: "docker-worker",
Backend: config.BackendConfig{
Docker: &config.DockerConfig{
SidecarImage: "my-registry.io/warpdotdev/warp-agent:latest",
},
},
}

wc, err := mergeConfig(fileConfig)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}

if wc.BackendType != "docker" {
t.Fatalf("BackendType = %q, want %q", wc.BackendType, "docker")
}
if wc.Docker == nil {
t.Fatal("expected docker backend config")
}
if wc.Docker.SidecarImage != "my-registry.io/warpdotdev/warp-agent:latest" {
t.Errorf("SidecarImage = %q, want %q", wc.Docker.SidecarImage, "my-registry.io/warpdotdev/warp-agent:latest")
}
}

func TestMergeConfigDockerSidecarImageOmitted(t *testing.T) {
resetCLIForTest()
t.Cleanup(resetCLIForTest)

fileConfig := &config.FileConfig{
WorkerID: "docker-worker",
Backend: config.BackendConfig{
Docker: &config.DockerConfig{},
},
}

wc, err := mergeConfig(fileConfig)
if err != nil {
t.Fatalf("unexpected error: %v", err)
}
if wc.Docker == nil {
t.Fatal("expected docker backend config")
}
if wc.Docker.SidecarImage != "" {
t.Errorf("SidecarImage = %q, want empty", wc.Docker.SidecarImage)
}
}

func TestMergeConfigDockerImagePullPolicyOmitted(t *testing.T) {
resetCLIForTest()
t.Cleanup(resetCLIForTest)
Expand Down
Loading