package docker import ( "context" "encoding/binary" "encoding/json" "errors" "io" "net/netip" "strings" "testing" "time" containertypes "github.com/moby/moby/api/types/container" imagetypes "github.com/moby/moby/api/types/image" mounttypes "github.com/moby/moby/api/types/mount" networktypes "github.com/moby/moby/api/types/network" volumetypes "github.com/moby/moby/api/types/volume" dockerclient "github.com/moby/moby/client" "tangled.org/samatkins.net/dockscope/internal/domain" ) type fakeLister struct { items []containertypes.Summary listErr error statsBody io.ReadCloser statsErr error logsBody io.ReadCloser logsErr error inspectTty bool inspectErr error images []imagetypes.Summary imagesErr error volumes []volumetypes.Volume volumesErr error networks []networktypes.Summary networksErr error history []imagetypes.HistoryResponseItem historyErr error inspectFull *containertypes.InspectResponse stopCalledID string stopCalledOpts dockerclient.ContainerStopOptions stopErr error restartCalledID string restartErr error removeCalledID string removeCalledOpts dockerclient.ContainerRemoveOptions removeErr error } func (f *fakeLister) ContainerList(_ context.Context, _ dockerclient.ContainerListOptions) (dockerclient.ContainerListResult, error) { if f.listErr != nil { return dockerclient.ContainerListResult{}, f.listErr } return dockerclient.ContainerListResult{Items: f.items}, nil } func (f *fakeLister) ContainerStop(_ context.Context, container string, opts dockerclient.ContainerStopOptions) (dockerclient.ContainerStopResult, error) { f.stopCalledID = container f.stopCalledOpts = opts return dockerclient.ContainerStopResult{}, f.stopErr } func (f *fakeLister) ContainerRestart(_ context.Context, container string, _ dockerclient.ContainerRestartOptions) (dockerclient.ContainerRestartResult, error) { f.restartCalledID = container return dockerclient.ContainerRestartResult{}, f.restartErr } func (f *fakeLister) ContainerRemove(_ context.Context, container string, opts dockerclient.ContainerRemoveOptions) (dockerclient.ContainerRemoveResult, error) { f.removeCalledID = container f.removeCalledOpts = opts return dockerclient.ContainerRemoveResult{}, f.removeErr } func (f *fakeLister) ContainerStats(_ context.Context, _ string, _ dockerclient.ContainerStatsOptions) (dockerclient.ContainerStatsResult, error) { if f.statsErr != nil { return dockerclient.ContainerStatsResult{}, f.statsErr } return dockerclient.ContainerStatsResult{Body: f.statsBody}, nil } func (f *fakeLister) ContainerLogs(_ context.Context, _ string, _ dockerclient.ContainerLogsOptions) (dockerclient.ContainerLogsResult, error) { if f.logsErr != nil { return nil, f.logsErr } return f.logsBody, nil } func (f *fakeLister) ContainerInspect(_ context.Context, _ string, _ dockerclient.ContainerInspectOptions) (dockerclient.ContainerInspectResult, error) { if f.inspectErr != nil { return dockerclient.ContainerInspectResult{}, f.inspectErr } if f.inspectFull != nil { return dockerclient.ContainerInspectResult{Container: *f.inspectFull}, nil } return dockerclient.ContainerInspectResult{ Container: containertypes.InspectResponse{ Config: &containertypes.Config{Tty: f.inspectTty}, }, }, nil } func (f *fakeLister) ImageList(_ context.Context, _ dockerclient.ImageListOptions) (dockerclient.ImageListResult, error) { if f.imagesErr != nil { return dockerclient.ImageListResult{}, f.imagesErr } return dockerclient.ImageListResult{Items: f.images}, nil } func (f *fakeLister) VolumeList(_ context.Context, _ dockerclient.VolumeListOptions) (dockerclient.VolumeListResult, error) { if f.volumesErr != nil { return dockerclient.VolumeListResult{}, f.volumesErr } return dockerclient.VolumeListResult{Items: f.volumes}, nil } func (f *fakeLister) NetworkList(_ context.Context, _ dockerclient.NetworkListOptions) (dockerclient.NetworkListResult, error) { if f.networksErr != nil { return dockerclient.NetworkListResult{}, f.networksErr } return dockerclient.NetworkListResult{Items: f.networks}, nil } func (f *fakeLister) ImageHistory(_ context.Context, _ string, _ ...dockerclient.ImageHistoryOption) (dockerclient.ImageHistoryResult, error) { if f.historyErr != nil { return dockerclient.ImageHistoryResult{}, f.historyErr } return dockerclient.ImageHistoryResult{Items: f.history}, nil } func (f *fakeLister) Close() error { return nil } // muxFrame builds a single Docker-style multiplexed log frame. stream is // 1 (stdout) or 2 (stderr). func muxFrame(stream byte, payload string) []byte { hdr := make([]byte, 8) hdr[0] = stream binary.BigEndian.PutUint32(hdr[4:8], uint32(len(payload))) return append(hdr, []byte(payload)...) } // statsBodyFromRecords marshals the given records to a stream of // concatenated JSON objects (the wire format the daemon emits when // Stream=true is set on ContainerStats). func statsBodyFromRecords(t *testing.T, records ...containertypes.StatsResponse) io.ReadCloser { t.Helper() var b strings.Builder for _, r := range records { j, err := json.Marshal(r) if err != nil { t.Fatalf("marshal stats: %v", err) } b.Write(j) b.WriteString("\n") } return io.NopCloser(strings.NewReader(b.String())) } func TestListContainersStripsLeadingSlashFromName(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ {ID: strings.Repeat("a", 64), Names: []string{"/api-gateway"}, State: "running"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } if len(got) != 1 { t.Fatalf("got %d containers, want 1", len(got)) } if got[0].Name != "api-gateway" { t.Errorf("Name = %q, want %q", got[0].Name, "api-gateway") } } func TestListContainersTruncatesLongID(t *testing.T) { t.Parallel() fullID := strings.Repeat("b", 64) cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ {ID: fullID, Names: []string{"/svc"}, State: "running"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } if got[0].ID != fullID[:shortIDLen] { t.Errorf("ID = %q, want %q", got[0].ID, fullID[:shortIDLen]) } } func TestListContainersHandlesShortID(t *testing.T) { t.Parallel() // Defends against a panic if the daemon ever returns an ID shorter than shortIDLen. cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ {ID: "abc", Names: []string{"/svc"}, State: "running"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } if got[0].ID != "abc" { t.Errorf("ID = %q, want %q", got[0].ID, "abc") } } func TestListContainersHandlesEmptyNames(t *testing.T) { t.Parallel() // Defends against a panic if the daemon ever returns a container with no names. cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ {ID: strings.Repeat("c", 64), Names: nil, State: "running"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } if got[0].Name != "" { t.Errorf("Name = %q, want empty", got[0].Name) } } func TestListContainersMapsState(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ {ID: strings.Repeat("d", 64), Names: []string{"/a"}, State: "running"}, {ID: strings.Repeat("e", 64), Names: []string{"/b"}, State: "exited"}, {ID: strings.Repeat("f", 64), Names: []string{"/c"}, State: "weird-future-state"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } want := []domain.ContainerState{ domain.StateRunning, domain.StateExited, domain.StateUnknown, } for i, w := range want { if got[i].State != w { t.Errorf("container %d: State = %q, want %q", i, got[i].State, w) } } } func TestListContainersMapsComposeLabels(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{ items: []containertypes.Summary{ { ID: "abc123", Names: []string{"/shop-web-1"}, State: "running", Labels: map[string]string{ "com.docker.compose.project": "shop", "com.docker.compose.service": "web", }, }, {ID: "def456", Names: []string{"/redis"}, State: "running"}, }, }} got, err := cli.ListContainers(t.Context()) if err != nil { t.Fatalf("ListContainers: %v", err) } if got[0].Project != "shop" || got[0].Service != "web" { t.Errorf("got Project=%q Service=%q, want shop/web", got[0].Project, got[0].Service) } if got[1].Project != "" || got[1].Service != "" { t.Errorf("unlabelled container got Project=%q Service=%q, want empty", got[1].Project, got[1].Service) } } func TestListImagesSkipsDanglingAndTrimsID(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{images: []imagetypes.Summary{ {ID: "sha256:" + strings.Repeat("a", 64), RepoTags: []string{"web:latest"}, Size: 100, Created: 1234}, {ID: "sha256:" + strings.Repeat("b", 64), RepoTags: []string{":"}}, {ID: "sha256:" + strings.Repeat("c", 64)}, }}} got, err := cli.ListImages(t.Context()) if err != nil { t.Fatalf("ListImages: %v", err) } if len(got) != 1 { t.Fatalf("got %d images, want 1", len(got)) } if got[0].ID != strings.Repeat("a", 12) { t.Errorf("ID = %q, want 12-char trimmed", got[0].ID) } if got[0].RepoTags[0] != "web:latest" || got[0].Size != 100 || got[0].Created != 1234 { t.Errorf("unexpected image: %+v", got[0]) } } func TestListImagesWrapsError(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{imagesErr: errors.New("boom")}} _, err := cli.ListImages(t.Context()) if err == nil || !strings.Contains(err.Error(), "list images") { t.Errorf("err = %v, want wrapped 'list images' error", err) } } func TestListVolumesMapsFields(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{volumes: []volumetypes.Volume{ { Name: "pgdata", Driver: "local", Mountpoint: "/var/lib/docker/volumes/pgdata/_data", Scope: "local", CreatedAt: "2026-07-01T10:00:00Z", Labels: map[string]string{"com.docker.compose.project": "shop"}, Options: map[string]string{"type": "none"}, }, }}} got, err := cli.ListVolumes(t.Context()) if err != nil { t.Fatalf("ListVolumes: %v", err) } if len(got) != 1 { t.Fatalf("got %d volumes, want 1", len(got)) } v := got[0] if v.Name != "pgdata" || v.Driver != "local" || v.Mountpoint == "" || v.Scope != "local" { t.Errorf("unexpected volume: %+v", v) } if v.Labels["com.docker.compose.project"] != "shop" || v.Options["type"] != "none" { t.Errorf("labels/options not mapped: %+v", v) } } func TestListNetworksExtractsSubnetAndGateway(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{networks: []networktypes.Summary{ { Network: networktypes.Network{ Name: "shop_default", ID: strings.Repeat("n", 64), Driver: "bridge", Scope: "local", IPAM: networktypes.IPAM{Config: []networktypes.IPAMConfig{{ Subnet: netip.MustParsePrefix("172.18.0.0/16"), Gateway: netip.MustParseAddr("172.18.0.1"), }}}, }, }, {Network: networktypes.Network{Name: "bridge", ID: "short", Driver: "bridge"}}, }}} got, err := cli.ListNetworks(t.Context()) if err != nil { t.Fatalf("ListNetworks: %v", err) } if len(got) != 2 { t.Fatalf("got %d networks, want 2", len(got)) } if got[0].Subnet != "172.18.0.0/16" || got[0].Gateway != "172.18.0.1" { t.Errorf("got Subnet=%q Gateway=%q", got[0].Subnet, got[0].Gateway) } if got[0].ID != strings.Repeat("n", 12) { t.Errorf("ID = %q, want 12-char trimmed", got[0].ID) } if got[1].Subnet != "" || got[1].Gateway != "" { t.Errorf("empty IPAM got Subnet=%q Gateway=%q, want empty", got[1].Subnet, got[1].Gateway) } } func TestInspectContainerConfigMapsFields(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{inspectFull: &containertypes.InspectResponse{ Created: "2026-07-01T10:00:00Z", State: &containertypes.State{Status: "running"}, Config: &containertypes.Config{ Env: []string{"B=2", "A=1"}, Image: "web:latest", Cmd: []string{"run", "--dev"}, Entrypoint: []string{"/bin/sh", "-c"}, WorkingDir: "/app", User: "app", Labels: map[string]string{"team": "core"}, ExposedPorts: networktypes.PortSet{ networktypes.MustParsePort("8080/tcp"): {}, }, }, Mounts: []containertypes.MountPoint{ {Type: mounttypes.TypeVolume, Name: "pgdata", Source: "/host/path", Destination: "/data", RW: true}, }, NetworkSettings: &containertypes.NetworkSettings{ Networks: map[string]*networktypes.EndpointSettings{ "bridge": { IPAddress: netip.MustParseAddr("172.17.0.2"), Gateway: netip.MustParseAddr("172.17.0.1"), }, }, }, }}} got, err := cli.InspectContainerConfig(t.Context(), "abc123") if err != nil { t.Fatalf("InspectContainerConfig: %v", err) } if got.Image != "web:latest" || got.Status != "running" || got.Created == "" { t.Errorf("unexpected base fields: %+v", got) } if len(got.Env) != 2 || got.Env[0] != "B=2" { t.Errorf("Env not passed through: %v", got.Env) } if len(got.ExposedPorts) != 1 || got.ExposedPorts[0] != "8080/tcp" { t.Errorf("ExposedPorts = %v", got.ExposedPorts) } if len(got.Mounts) != 1 || got.Mounts[0].Name != "pgdata" || !got.Mounts[0].RW { t.Errorf("Mounts = %+v", got.Mounts) } if len(got.Networks) != 1 || got.Networks[0].IPAddress != "172.17.0.2" || got.Networks[0].Gateway != "172.17.0.1" { t.Errorf("Networks = %+v", got.Networks) } if got.Labels["team"] != "core" { t.Errorf("Labels not mapped: %v", got.Labels) } } func TestImageHistoryMapsLayers(t *testing.T) { t.Parallel() cli := &liveClient{cli: &fakeLister{history: []imagetypes.HistoryResponseItem{ {ID: "sha256:" + strings.Repeat("a", 64), Size: 50, CreatedBy: "RUN make build", Tags: []string{"web:latest"}}, {ID: "", Size: 0, CreatedBy: `CMD ["sh"]`}, }}} got, err := cli.ImageHistory(t.Context(), "abc123") if err != nil { t.Fatalf("ImageHistory: %v", err) } if len(got) != 2 { t.Fatalf("got %d layers, want 2", len(got)) } if got[0].ID != strings.Repeat("a", 12) || got[0].Size != 50 || got[0].Command != "RUN make build" { t.Errorf("unexpected layer: %+v", got[0]) } if got[1].ID != "" { t.Errorf("missing layer ID = %q, want ", got[1].ID) } } func TestListContainersPropagatesError(t *testing.T) { t.Parallel() wantErr := errors.New("docker daemon unreachable") cli := &liveClient{cli: &fakeLister{listErr: wantErr}} _, err := cli.ListContainers(t.Context()) if err == nil { t.Fatal("expected error, got nil") } if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } func TestStreamStatsReceivesStats(t *testing.T) { t.Parallel() record := containertypes.StatsResponse{ CPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 1000}, SystemUsage: 5000, OnlineCPUs: 4, }, MemoryStats: containertypes.MemoryStats{Usage: 1048576, Limit: 2097152}, PidsStats: containertypes.PidsStats{Current: 42}, } cli := &liveClient{cli: &fakeLister{ statsBody: statsBodyFromRecords(t, record), }} ch, err := cli.StreamStats(t.Context(), "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } select { case stats := <-ch: if stats.MemUsage != 1048576 { t.Errorf("MemUsage = %d, want 1048576", stats.MemUsage) } if stats.MemLimit != 2097152 { t.Errorf("MemLimit = %d, want 2097152", stats.MemLimit) } if stats.PIDs != 42 { t.Errorf("PIDs = %d, want 42", stats.PIDs) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for stats") } } func TestStreamStatsCPUPercentFromPreCPUStats(t *testing.T) { t.Parallel() // First record: PreRead is zero → CPUPercent should be 0 (no previous sample). first := containertypes.StatsResponse{ CPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 1000}, SystemUsage: 5000, OnlineCPUs: 2, }, } // Second record: PreRead is set, daemon supplies PreCPUStats. We compute // (cpuDelta / systemDelta) * numCPUs * 100. second := containertypes.StatsResponse{ PreRead: time.Now(), PreCPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 1000}, SystemUsage: 5000, }, CPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 3000}, SystemUsage: 10000, OnlineCPUs: 2, }, } cli := &liveClient{cli: &fakeLister{statsBody: statsBodyFromRecords(t, first, second)}} ch, err := cli.StreamStats(t.Context(), "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } select { case stats := <-ch: if stats.CPUPercent != 0 { t.Errorf("first CPUPercent = %f, want 0 (no previous sample)", stats.CPUPercent) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for first stats") } select { case stats := <-ch: want := (float64(2000) / float64(5000)) * 2 * 100.0 if stats.CPUPercent != want { t.Errorf("second CPUPercent = %f, want %f", stats.CPUPercent, want) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for second stats") } } func TestStreamStatsNetworkRates(t *testing.T) { t.Parallel() t0 := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) first := containertypes.StatsResponse{ Read: t0, Networks: map[string]containertypes.NetworkStats{ "eth0": {RxBytes: 1000, TxBytes: 500}, }, } second := containertypes.StatsResponse{ Read: t0.Add(2 * time.Second), Networks: map[string]containertypes.NetworkStats{ "eth0": {RxBytes: 5000, TxBytes: 2500}, }, } cli := &liveClient{cli: &fakeLister{statsBody: statsBodyFromRecords(t, first, second)}} ch, err := cli.StreamStats(t.Context(), "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } select { case stats := <-ch: if stats.NetRxBps != 0 || stats.NetTxBps != 0 { t.Errorf("first sample rates = (%f, %f), want (0, 0)", stats.NetRxBps, stats.NetTxBps) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for first stats") } select { case stats := <-ch: wantRx := float64(4000) / 2 wantTx := float64(2000) / 2 if stats.NetRxBps != wantRx { t.Errorf("NetRxBps = %f, want %f", stats.NetRxBps, wantRx) } if stats.NetTxBps != wantTx { t.Errorf("NetTxBps = %f, want %f", stats.NetTxBps, wantTx) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for second stats") } } func TestStreamStatsNetworkCounterReset(t *testing.T) { t.Parallel() t0 := time.Date(2026, 1, 1, 0, 0, 0, 0, time.UTC) first := containertypes.StatsResponse{ Read: t0, Networks: map[string]containertypes.NetworkStats{ "eth0": {RxBytes: 5000, TxBytes: 5000}, }, } second := containertypes.StatsResponse{ Read: t0.Add(1 * time.Second), Networks: map[string]containertypes.NetworkStats{ "eth0": {RxBytes: 100, TxBytes: 100}, }, } cli := &liveClient{cli: &fakeLister{statsBody: statsBodyFromRecords(t, first, second)}} ch, err := cli.StreamStats(t.Context(), "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } <-ch // skip first select { case stats := <-ch: if stats.NetRxBps != 0 || stats.NetTxBps != 0 { t.Errorf("counter reset rates = (%f, %f), want (0, 0)", stats.NetRxBps, stats.NetTxBps) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for second stats") } } // TestStreamStatsCPUPercentCgroupV2 covers the cgroup-v2 path where // PercpuUsage is empty but OnlineCPUs is populated. This is the path real // users hit on modern Linux daemons. func TestStreamStatsCPUPercentCgroupV2(t *testing.T) { t.Parallel() record := containertypes.StatsResponse{ PreRead: time.Now(), PreCPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 100, PercpuUsage: nil}, SystemUsage: 1000, }, CPUStats: containertypes.CPUStats{ CPUUsage: containertypes.CPUUsage{TotalUsage: 200, PercpuUsage: nil}, SystemUsage: 2000, OnlineCPUs: 8, }, } cli := &liveClient{cli: &fakeLister{statsBody: statsBodyFromRecords(t, record)}} ch, err := cli.StreamStats(t.Context(), "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } select { case stats := <-ch: // (100/1000) * 8 * 100 = 80 want := 80.0 if stats.CPUPercent != want { t.Errorf("CPUPercent = %f, want %f (cgroup v2 should use OnlineCPUs)", stats.CPUPercent, want) } case <-time.After(2 * time.Second): t.Fatal("timeout waiting for stats") } } // TestStreamStatsFetchErrorReturnsErr verifies that a fetch-time error is // surfaced from StreamStats itself rather than swallowed. func TestStreamStatsFetchErrorReturnsErr(t *testing.T) { t.Parallel() wantErr := errors.New("stats unavailable") cli := &liveClient{cli: &fakeLister{statsErr: wantErr}} _, err := cli.StreamStats(t.Context(), "abc") if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } // TestStreamStatsClosesChannelOnContextCancel verifies the goroutine // terminates and closes the channel when its context is cancelled. func TestStreamStatsClosesChannelOnContextCancel(t *testing.T) { t.Parallel() // A reader that blocks until the context is cancelled, simulating a // long-lived stream. body := &blockingReader{done: make(chan struct{})} cli := &liveClient{cli: &fakeLister{statsBody: body}} ctx, cancel := context.WithCancel(t.Context()) ch, err := cli.StreamStats(ctx, "abc") if err != nil { t.Fatalf("StreamStats: %v", err) } cancel() close(body.done) // unblock Read so the decoder errors out select { case _, ok := <-ch: if ok { t.Error("channel produced a value after cancel; expected close") } case <-time.After(2 * time.Second): t.Fatal("channel was not closed after context cancel") } } // blockingReader returns 0 bytes / no error on Read until done is closed, // after which it returns io.EOF. type blockingReader struct { done chan struct{} } func (b *blockingReader) Read(p []byte) (int, error) { <-b.done return 0, io.EOF } func (b *blockingReader) Close() error { return nil } // TestStreamLogsDemuxesNonTtyContainer verifies that a non-TTY container's // 8-byte-framed stream is demultiplexed correctly. func TestStreamLogsDemuxesNonTtyContainer(t *testing.T) { t.Parallel() body := append(muxFrame(1, "hello\n"), muxFrame(2, "world\n")...) cli := &liveClient{cli: &fakeLister{ logsBody: io.NopCloser(strings.NewReader(string(body))), inspectTty: false, }} rc, err := cli.StreamLogs(t.Context(), "abc", "100") if err != nil { t.Fatalf("StreamLogs: %v", err) } defer func() { _ = rc.Close() }() got, err := io.ReadAll(rc) if err != nil { t.Fatalf("read: %v", err) } if string(got) != "hello\nworld\n" { t.Errorf("got %q, want %q", got, "hello\nworld\n") } } // TestStreamLogsPassesThroughTtyContainer verifies that a TTY container's // raw bytes are passed through without demultiplexing (which would corrupt them). func TestStreamLogsPassesThroughTtyContainer(t *testing.T) { t.Parallel() raw := "2026-05-09T20:06:09.000000000Z hello tty\n" cli := &liveClient{cli: &fakeLister{ logsBody: io.NopCloser(strings.NewReader(raw)), inspectTty: true, }} rc, err := cli.StreamLogs(t.Context(), "abc", "100") if err != nil { t.Fatalf("StreamLogs: %v", err) } defer func() { _ = rc.Close() }() got, err := io.ReadAll(rc) if err != nil { t.Fatalf("read: %v", err) } if string(got) != raw { t.Errorf("got %q, want %q", got, raw) } } // TestStreamLogsInspectErrorPropagates verifies that an inspect failure is // surfaced rather than swallowed. func TestStreamLogsInspectErrorPropagates(t *testing.T) { t.Parallel() wantErr := errors.New("inspect blew up") cli := &liveClient{cli: &fakeLister{inspectErr: wantErr}} _, err := cli.StreamLogs(t.Context(), "abc", "100") if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } func TestStopContainerErrorWraps(t *testing.T) { t.Parallel() wantErr := errors.New("stop failed") lister := &fakeLister{stopErr: wantErr} cli := &liveClient{cli: lister} err := cli.StopContainer(t.Context(), "abc") if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } func TestRestartContainerErrorWraps(t *testing.T) { t.Parallel() wantErr := errors.New("restart failed") lister := &fakeLister{restartErr: wantErr} cli := &liveClient{cli: lister} err := cli.RestartContainer(t.Context(), "abc") if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } func TestStopContainerForwardsArgs(t *testing.T) { t.Parallel() lister := &fakeLister{} cli := &liveClient{cli: lister} if err := cli.StopContainer(t.Context(), "abc"); err != nil { t.Fatalf("StopContainer: %v", err) } if lister.stopCalledID != "abc" { t.Errorf("forwarded ID = %q, want %q", lister.stopCalledID, "abc") } } func TestRestartContainerForwardsArgs(t *testing.T) { t.Parallel() lister := &fakeLister{} cli := &liveClient{cli: lister} if err := cli.RestartContainer(t.Context(), "abc"); err != nil { t.Fatalf("RestartContainer: %v", err) } if lister.restartCalledID != "abc" { t.Errorf("forwarded ID = %q, want %q", lister.restartCalledID, "abc") } } func TestRemoveContainerSuccess(t *testing.T) { t.Parallel() lister := &fakeLister{} cli := &liveClient{cli: lister} if err := cli.RemoveContainer(t.Context(), "abc"); err != nil { t.Fatalf("RemoveContainer: %v", err) } if lister.removeCalledID != "abc" { t.Errorf("forwarded ID = %q, want %q", lister.removeCalledID, "abc") } } func TestRemoveContainerErrorWraps(t *testing.T) { t.Parallel() wantErr := errors.New("remove failed") lister := &fakeLister{removeErr: wantErr} cli := &liveClient{cli: lister} err := cli.RemoveContainer(t.Context(), "abc") if !errors.Is(err, wantErr) { t.Errorf("err = %v, want wrap of %v", err, wantErr) } } func TestRemoveContainerPreservesVolumes(t *testing.T) { t.Parallel() lister := &fakeLister{} cli := &liveClient{cli: lister} if err := cli.RemoveContainer(t.Context(), "abc"); err != nil { t.Fatalf("RemoveContainer: %v", err) } if !lister.removeCalledOpts.Force { t.Errorf("Force = false, want true") } if lister.removeCalledOpts.RemoveVolumes { t.Errorf("RemoveVolumes = true, want false (named volumes must survive)") } }