Skip to content
Open
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
43 changes: 8 additions & 35 deletions pkg/collector/collector.go
Original file line number Diff line number Diff line change
@@ -1,18 +1,15 @@
package main

import (
"encoding/json"
"flag"
"fmt"
"io"
"net/http"
_ "net/http/pprof"
"os"
"time"

"github.com/eraser-dev/eraser/pkg/cri"
"github.com/eraser-dev/eraser/pkg/logger"
"golang.org/x/sys/unix"
logf "sigs.k8s.io/controller-runtime/pkg/log"

util "github.com/eraser-dev/eraser/pkg/utils"
Expand Down Expand Up @@ -73,57 +70,33 @@ func main() {
}
log.Info("images collected", "finalImages:", finalImages)

data, err := json.Marshal(finalImages)
if err != nil {
log.Error(err, "failed to encode finalImages")
os.Exit(1)
}

path := util.CollectScanPath

if *scanDisabled {
path = util.ScanErasePath
}

if err := unix.Mkfifo(path, util.PipeMode); err != nil {
log.Error(err, "failed to create pipe", "pipeFile", path)
os.Exit(1)
}

//nolint:gosec // G304: Opening pipe file is intended functionality
file, err := os.OpenFile(path, os.O_WRONLY, 0)
// Published before the payload, not after: the peer can finish and signal
// back the moment it has read the list, so an endpoint created afterwards
// can be missed entirely. The scanner already publishes in this order.
completion, err := util.CreateCompletionPipe(util.EraseCompleteCollectPath)
if err != nil {
log.Error(err, "failed to open pipe", "pipeFile", path)
os.Exit(1)
}

if _, err := file.Write(data); err != nil {
log.Error(err, "failed to write to pipe", "pipeFile", path)
os.Exit(1)
}

if err := file.Close(); err != nil {
log.Error(err, "failed to close pipe", "pipeFile", path)
os.Exit(1)
}
if err := unix.Mkfifo(util.EraseCompleteCollectPath, util.PipeMode); err != nil {
log.Error(err, "failed to create pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(1)
}
Comment thread
charleswool marked this conversation as resolved.

file, err = os.OpenFile(util.EraseCompleteCollectPath, os.O_RDONLY, 0)
if err != nil {
log.Error(err, "failed to open pipe", "pipeFile", util.EraseCompleteCollectPath)
if err := util.WriteImagesPipe(path, finalImages); err != nil {
log.Error(err, "failed to send images", "pipeFile", path)
os.Exit(1)
}

data, err = io.ReadAll(file)
data, err := completion.Await()
if err != nil {
log.Error(err, "failed to read pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(1)
}

if err := file.Close(); err != nil {
if err := completion.Close(); err != nil {
log.Error(err, "failed to close pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(1)
}
Expand Down
62 changes: 5 additions & 57 deletions pkg/remover/remover.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,11 +2,8 @@ package main

import (
"context"
"encoding/json"
"flag"
"fmt"
"io"
"io/fs"
"net/http"
_ "net/http/pprof"
"os"
Expand All @@ -21,7 +18,6 @@ import (
"github.com/eraser-dev/eraser/pkg/logger"
"github.com/eraser-dev/eraser/pkg/metrics"

"github.com/eraser-dev/eraser/api/unversioned"
util "github.com/eraser-dev/eraser/pkg/utils"
)

Expand Down Expand Up @@ -78,38 +74,11 @@ func main() {
}

if *imageListPtr == "" {
var f *os.File
for {
var err error

f, err = os.OpenFile(util.ScanErasePath, os.O_RDONLY, 0)
if err == nil {
break
}
if !os.IsNotExist(err) {
log.Error(err, "error opening scanErase pipe")
os.Exit(generalErr)
}
time.Sleep(1 * time.Second)
continue
}

// json data is list of []unversioned.Image
data, err := io.ReadAll(f)
nonCompliantImages, err := util.ReadImagesPipe(context.Background(), util.ScanErasePath)
if err != nil {
log.Error(err, "error reading non-compliant images")
os.Exit(generalErr)
}
if err := f.Close(); err != nil {
log.Error(err, "error closing non-compliant images file")
os.Exit(generalErr)
}

nonCompliantImages := []unversioned.Image{}
if err = json.Unmarshal(data, &nonCompliantImages); err != nil {
log.Error(err, "error in unmarshal non-compliant images")
os.Exit(generalErr)
}

for _, img := range nonCompliantImages {
imagelist = append(imagelist, img.ImageID)
Expand Down Expand Up @@ -158,39 +127,18 @@ func main() {
}

if *imageListPtr == "" {
file, err := os.OpenFile(util.EraseCompleteCollectPath, os.O_WRONLY, 0)
if err != nil {
log.Error(err, "unable to open pipe", "pipeFile", util.EraseCompleteCollectPath)
if err := util.WriteCompletionPipe(util.EraseCompleteCollectPath); err != nil {
log.Error(err, "unable to signal completion", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(generalErr)
}

if _, err := file.WriteString(util.EraseCompleteMessage); err != nil {
log.Error(err, "unable to write to pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(generalErr)
}

if err := file.Close(); err != nil {
log.Error(err, "unable to close pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(generalErr)
}

file, err = os.OpenFile(util.EraseCompleteScanPath, os.O_WRONLY, fs.ModeNamedPipe)
err := util.WriteCompletionPipe(util.EraseCompleteScanPath)
// if the scanner is disabled
if os.IsNotExist(err) {
return
}
if err != nil {
log.Error(err, "unable to open pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(generalErr)
}

if _, err := file.WriteString(util.EraseCompleteMessage); err != nil {
log.Error(err, "unable to write to pipe", "pipeFile", util.EraseCompleteCollectPath)
os.Exit(generalErr)
}

if err := file.Close(); err != nil {
log.Error(err, "unable to close pipe", "pipeFile", util.EraseCompleteScanPath)
log.Error(err, "unable to signal completion", "pipeFile", util.EraseCompleteScanPath)
os.Exit(generalErr)
}
}
Expand Down
23 changes: 9 additions & 14 deletions pkg/scanners/template/scanner_template.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,14 +2,12 @@ package template

import (
"context"
"io"
"os"
"os/signal"
"syscall"

"github.com/eraser-dev/eraser/api/unversioned"
"github.com/go-logr/logr"
"golang.org/x/sys/unix"

"github.com/eraser-dev/eraser/pkg/metrics"
util "github.com/eraser-dev/eraser/pkg/utils"
Expand All @@ -35,6 +33,10 @@ type config struct {
deleteScanFailedImages bool
deleteEOLImages bool
reportMetrics bool

// held from ReceiveImages until Finish: the endpoint must stay published so
// the remover can tell a scanner is present
completion *util.CompletionPipe
}

type ConfigFunc func(*config)
Expand All @@ -59,7 +61,9 @@ func NewImageProvider(funcs ...ConfigFunc) ImageProvider {
func (cfg *config) ReceiveImages() ([]unversioned.Image, error) {
var err error

if err := unix.Mkfifo(util.EraseCompleteScanPath, util.PipeMode); err != nil {
// published up front so the remover can tell a scanner is present
cfg.completion, err = util.CreateCompletionPipe(util.EraseCompleteScanPath)
if err != nil {
cfg.log.Error(err, "failed to create pipe", "pipeName", util.EraseCompleteScanPath)
return nil, err
}
Expand Down Expand Up @@ -107,23 +111,14 @@ func (cfg *config) SendImages(nonCompliantImages, failedImages []unversioned.Ima
}

func (cfg *config) Finish() error {
file, err := os.OpenFile(util.EraseCompleteScanPath, os.O_RDONLY, 0)
if err != nil {
cfg.log.Error(err, "failed to open pipe", "pipeName", util.EraseCompleteScanPath)
return err
}
defer func() { _ = cfg.completion.Close() }()

data, err := io.ReadAll(file)
data, err := cfg.completion.Await()
if err != nil {
cfg.log.Error(err, "failed to read pipe", "pipeName", util.EraseCompleteScanPath)
return err
}

if err := file.Close(); err != nil {
cfg.log.Error(err, "failed to close pipe", "pipeName", util.EraseCompleteScanPath)
return err
}

if string(data) != util.EraseCompleteMessage {
cfg.log.Info("garbage in pipe", "pipeName", util.EraseCompleteScanPath, "in_pipe", string(data))
return err
Expand Down
113 changes: 113 additions & 0 deletions pkg/utils/handoff_test.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,113 @@
package utils

import (
"context"
"os"
"path/filepath"
"testing"

"github.com/eraser-dev/eraser/api/unversioned"
)

// These run against whichever implementation the platform selects: FIFOs on
// Unix, Unix domain sockets on Windows. Keeping them build-tag free is the point
// -- the two implementations have to stay behaviorally identical.

// shortTempDir keeps paths well inside the sun_path limit that applies to the
// socket implementation.
func shortTempDir(t *testing.T) string {
t.Helper()

dir, err := os.MkdirTemp("", "h")
if err != nil {
t.Fatal(err)
}
t.Cleanup(func() { _ = os.RemoveAll(dir) })

return dir
}

func TestImagesHandoffRoundTrip(t *testing.T) {
path := filepath.Join(shortTempDir(t), "images")

want := []unversioned.Image{
{ImageID: "sha256:aaaa", Names: []string{"repo/one:v1"}},
{ImageID: "sha256:bbbb", Names: []string{"repo/two:v2"}},
}

errCh := make(chan error, 1)
go func() { errCh <- WriteImagesPipe(path, want) }()

got, err := ReadImagesPipe(context.Background(), path)
if err != nil {
t.Fatalf("ReadImagesPipe: %v", err)
}
if err := <-errCh; err != nil {
t.Fatalf("WriteImagesPipe: %v", err)
}

if len(got) != len(want) {
t.Fatalf("got %d images, want %d", len(got), len(want))
}
for i := range want {
if got[i].ImageID != want[i].ImageID {
t.Errorf("image %d = %q, want %q", i, got[i].ImageID, want[i].ImageID)
}
}
}

func TestCompletionHandoffRoundTrip(t *testing.T) {
path := filepath.Join(shortTempDir(t), "complete")

pipe, err := CreateCompletionPipe(path)
if err != nil {
t.Fatalf("CreateCompletionPipe: %v", err)
}
defer func() { _ = pipe.Close() }()

errCh := make(chan error, 1)
go func() { errCh <- WriteCompletionPipe(path) }()

data, err := pipe.Await()
if err != nil {
t.Fatalf("Await: %v", err)
}
if err := <-errCh; err != nil {
t.Fatalf("WriteCompletionPipe: %v", err)
}

if string(data) != EraseCompleteMessage {
t.Errorf("payload = %q, want %q", string(data), EraseCompleteMessage)
}
}

// The remover infers "the scanner is disabled" from this error, so it has to
// keep satisfying os.IsNotExist on both platforms.
func TestWriteCompletionPipeAbsentPeerIsNotExist(t *testing.T) {
path := filepath.Join(shortTempDir(t), "no-such-peer")

err := WriteCompletionPipe(path)
if err == nil {
t.Fatal("expected an error writing to an endpoint nobody published")
}
if !os.IsNotExist(err) {
t.Errorf("os.IsNotExist(%v) = false, want true", err)
}
}

func TestCompletionPipeCloseIsIdempotentlySafe(t *testing.T) {
path := filepath.Join(shortTempDir(t), "closed")

pipe, err := CreateCompletionPipe(path)
if err != nil {
t.Fatalf("CreateCompletionPipe: %v", err)
}

// the scanner defers Close and the collector also closes explicitly
if err := pipe.Close(); err != nil {
t.Errorf("first Close: %v", err)
}
if err := pipe.Close(); err != nil {
t.Errorf("second Close: %v", err)
}
}
Loading
Loading