diff --git a/.github/workflows/pipeline.yml b/.github/workflows/pipeline.yml new file mode 100644 index 0000000..5c2bd65 --- /dev/null +++ b/.github/workflows/pipeline.yml @@ -0,0 +1,63 @@ +name: Lint, test, and build + +on: + pull_request: + push: + branches: + - main + - master + +permissions: + contents: read + +jobs: + lint: + name: Lint + runs-on: ubuntu-latest + steps: + - name: Checkout code + uses: actions/checkout@v4 + + - name: Set up Go + uses: actions/setup-go@v5 + with: + go-version-file: go.mod + cache: true + + - name: Format + run: make fmt + + - name: Vet + run: make vet + + test: + name: Tests + runs-on: ubuntu-latest + steps: + - name: Checkout code + uses: actions/checkout@v4 + + - name: Set up Go + uses: actions/setup-go@v5 + with: + go-version-file: go.mod + cache: true + + - name: Run tests + run: make test + + build: + name: Build + runs-on: ubuntu-latest + steps: + - name: Checkout code + uses: actions/checkout@v4 + + - name: Set up Go + uses: actions/setup-go@v5 + with: + go-version-file: go.mod + cache: true + + - name: Build release binary + run: make build-release diff --git a/.gitignore b/.gitignore index 5f60593..e619aef 100644 --- a/.gitignore +++ b/.gitignore @@ -11,4 +11,7 @@ # Output of the go coverage tool, specifically when used with LiteIDE *.out +# Local build artifacts +bin/ + cmd/testing* diff --git a/Makefile b/Makefile new file mode 100644 index 0000000..c5b4e87 --- /dev/null +++ b/Makefile @@ -0,0 +1,93 @@ +export GO111MODULE=on +BUILD_DIR ?= bin +BUILD_TAGS ?= osusergo,netgo,sqlite_omit_load_extension +LDFLAGS ?= -extldflags '-static' -s -w +PKG ?= ./cmd/coriolis-logger + +.PHONY: all +all: build-release + +.PHONY: fmt +fmt: ## Format the code. + go fmt ./... + +.PHONY: vet +vet: ## Run static code analysis. + go vet ./... + +COVER_OUTDIR ?= $(BUILD_DIR) +COVER_OUTFILE_RAW ?= $(COVER_OUTDIR)/coverage.raw +COVER_OUTFILE_HTML ?= $(COVER_OUTDIR)/coverage.html +# Filter executed tests by the "TEST_RE" regex +TEST_RE ?= .* +# Skip tests that match the "TEST_SKIP_RE" regex +TEST_SKIP_RE ?= +TEST_CMD = go test ./... \ + -run "$(TEST_RE)" \ + -coverpkg=coriolis-logger/... \ + -coverprofile=$(COVER_OUTFILE_RAW) +ifneq ($(strip $(TEST_SKIP_RE)),) +TEST_CMD += -skip "$(TEST_SKIP_RE)" +endif + +.PHONY: test-unit +test-unit: fmt vet ## Run coriolis-logger unit tests. + mkdir -p $(COVER_OUTDIR) + $(TEST_CMD) + go tool cover -html=$(COVER_OUTFILE_RAW) -o=$(COVER_OUTFILE_HTML) + +.PHONY: test-unit-verbose +test-unit-verbose: fmt vet ## Run coriolis-logger unit tests in verbose mode. + mkdir -p $(COVER_OUTDIR) + $(TEST_CMD) -test.v + go tool cover -html=$(COVER_OUTFILE_RAW) -o=$(COVER_OUTFILE_HTML) + +.PHONY: test +test: test-unit ## Run all coriolis-logger tests. + +.PHONY: build-dev +build-dev: fmt vet ## Generate coriolis-logger dev build. + # Dev build, meant to build fast and run on the dev machine: + # * use the host architecture + # * avoid rebuilding unmodified components + mkdir -p $(BUILD_DIR) + go build \ + -o $(BUILD_DIR)/coriolis-logger \ + $(PKG) + +.PHONY: build-dev-dbg +build-dev-dbg: fmt vet ## Generate coriolis-logger dev build, disabling compiler optimizations. + mkdir -p $(BUILD_DIR) + go build -gcflags="all=-N -l" \ + -o $(BUILD_DIR)/coriolis-logger \ + $(PKG) + +.PHONY: build-release +build-release: fmt vet ## Generate coriolis-logger release build. + # Release build, meant to be deployed as the Coriolis logging service: + # * strip debug symbols + # * build for Linux x86_64 + # * statically linked + # * rebuild everything + mkdir -p $(BUILD_DIR) + CGO_ENABLED=0 GOOS=linux GOARCH=amd64 go build -a \ + -tags "$(BUILD_TAGS)" \ + -ldflags="$(LDFLAGS)" \ + -o $(BUILD_DIR)/coriolis-logger \ + $(PKG) + +CONFIG ?= testdata/config.toml + +.PHONY: run +run: ## Run coriolis-logger. + $(BUILD_DIR)/coriolis-logger -config "$(CONFIG)" + +# We'll reuse the "help" generator from operator-sdk (Apache-2). +.DEFAULT_GOAL := help +.PHONY: help +help: ## Show this help screen. + @echo 'Usage: make ... ' + @echo '' + @echo 'Available targets are:' + @echo '' + @awk 'BEGIN {FS = ":.*##"; printf "\nUsage:\n make \033[36m\033[0m\n"} /^[a-zA-Z0-9_-]+:.*?##/ { printf " \033[36m%-25s\033[0m %s\n", $$1, $$2 } /^##@/ { printf "\n\033[1m%s\033[0m\n", substr($$0, 5) } ' $(MAKEFILE_LIST) diff --git a/README.md b/README.md index 188ff53..8e03701 100644 --- a/README.md +++ b/README.md @@ -22,6 +22,14 @@ Build the binary: ```bash cd coriolis-logger +make build-dev +``` + +This writes `bin/coriolis-logger`. Use `make build-release` for a Linux amd64 binary, and `make help` to list all targets. + +Alternatively: + +```bash go install ./... ``` @@ -63,6 +71,12 @@ listener = "unixgram" # address = "/tmp/coriolis-logger/syslog" address = "/tmp/coriolis-logging.sock" +# Optional second listener. The syslog server can bind a unix +# datagram socket and a TCP or UDP endpoint at the same time. +# extra_listener and extra_address must both be set when used. +# extra_listener = "tcp" +# extra_address = "0.0.0.0:5144" + # Log format # possible values: # rfc3164 diff --git a/cmd/coriolis-logger/main.go b/cmd/coriolis-logger/main.go index 690bea7..ba78429 100644 --- a/cmd/coriolis-logger/main.go +++ b/cmd/coriolis-logger/main.go @@ -35,7 +35,7 @@ import ( var log = loggo.GetLogger("coriolis.logger.cmd") func main() { - stop := make(chan os.Signal) + stop := make(chan os.Signal, 1) signal.Notify(stop, syscall.SIGTERM) signal.Notify(stop, syscall.SIGINT) log.SetLogLevel(loggo.DEBUG) diff --git a/config/config.go b/config/config.go index add2857..52f585c 100644 --- a/config/config.go +++ b/config/config.go @@ -176,12 +176,17 @@ func (a *APIServer) Validate() error { } type Syslog struct { - Listener ListenerType - Address string - Format string - LogToStdout bool `toml:"log_to_stdout"` - DataStore DatastoreType - InfluxDB *InfluxDB `toml:"influxdb"` + Listener ListenerType + Address string + // ExtraListener and ExtraAddress optionally start a second syslog + // endpoint on the same server. This allows a unix datagram socket + // to run alongside a TCP or UDP listener (or vice versa). + ExtraListener ListenerType `toml:"extra_listener"` + ExtraAddress string `toml:"extra_address"` + Format string + LogToStdout bool `toml:"log_to_stdout"` + DataStore DatastoreType + InfluxDB *InfluxDB `toml:"influxdb"` } func (s *Syslog) LogFormat() (format.Format, error) { @@ -213,9 +218,29 @@ func (s *Syslog) Validate() error { return fmt.Errorf("invalid datastore type %q", s.DataStore) } - switch s.Listener { + if err := validateListener(s.Listener, s.Address); err != nil { + return err + } + + if s.ExtraListener == "" && s.ExtraAddress == "" { + return nil + } + if s.ExtraListener == "" || s.ExtraAddress == "" { + return fmt.Errorf("extra_listener and extra_address must both be set") + } + if err := validateListener(s.ExtraListener, s.ExtraAddress); err != nil { + return errors.Wrap(err, "validating extra listener") + } + if s.Listener == s.ExtraListener && s.Address == s.ExtraAddress { + return fmt.Errorf("extra listener duplicates the primary listener") + } + return nil +} + +func validateListener(listener ListenerType, address string) error { + switch listener { case UnixDgramListener: - absPath, err := filepath.Abs(s.Address) + absPath, err := filepath.Abs(address) if err != nil { return errors.Wrap(err, "getting dirname") } @@ -224,15 +249,15 @@ func (s *Syslog) Validate() error { return errors.Wrap(err, "fetching info about dirname") } - if mode, err := os.Stat(s.Address); err == nil { + if mode, err := os.Stat(address); err == nil { if mode.Mode()&os.ModeSocket == 0 { return fmt.Errorf( - "cannot use %q as address. File already exists and is not socket", s.Address) + "cannot use %q as address. File already exists and is not socket", address) } } case TCPListener, UDPListener: default: - return fmt.Errorf("invalid listener type %q", s.Listener) + return fmt.Errorf("invalid listener type %q", listener) } return nil } diff --git a/config/config_test.go b/config/config_test.go new file mode 100644 index 0000000..3259b01 --- /dev/null +++ b/config/config_test.go @@ -0,0 +1,106 @@ +// Copyright 2026 Cloudbase Solutions SRL +// +// Licensed under the Apache License, Version 2.0 (the "License"); you may +// not use this file except in compliance with the License. You may obtain +// a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +// License for the specific language governing permissions and limitations +// under the License. + +package config + +import ( + "path/filepath" + "testing" +) + +func validSyslog(t *testing.T) Syslog { + t.Helper() + return Syslog{ + Listener: UnixDgramListener, + Address: filepath.Join(t.TempDir(), "syslog.sock"), + Format: "automatic", + DataStore: StdOutDataStore, + } +} + +func TestSyslogValidateExtraListener(t *testing.T) { + t.Parallel() + + t.Run("primary only", func(t *testing.T) { + cfg := validSyslog(t) + if err := cfg.Validate(); err != nil { + t.Fatalf("expected valid config, got %v", err) + } + }) + + t.Run("unix and tcp", func(t *testing.T) { + cfg := validSyslog(t) + cfg.ExtraListener = TCPListener + cfg.ExtraAddress = "127.0.0.1:5144" + if err := cfg.Validate(); err != nil { + t.Fatalf("expected valid unix+tcp config, got %v", err) + } + }) + + t.Run("unix and udp", func(t *testing.T) { + cfg := validSyslog(t) + cfg.ExtraListener = UDPListener + cfg.ExtraAddress = "127.0.0.1:5144" + if err := cfg.Validate(); err != nil { + t.Fatalf("expected valid unix+udp config, got %v", err) + } + }) + + t.Run("tcp primary with unix extra", func(t *testing.T) { + cfg := validSyslog(t) + cfg.Listener = TCPListener + cfg.Address = "0.0.0.0:5144" + cfg.ExtraListener = UnixDgramListener + cfg.ExtraAddress = filepath.Join(t.TempDir(), "extra.sock") + if err := cfg.Validate(); err != nil { + t.Fatalf("expected valid tcp+unix config, got %v", err) + } + }) + + t.Run("extra listener without address", func(t *testing.T) { + cfg := validSyslog(t) + cfg.ExtraListener = TCPListener + if err := cfg.Validate(); err == nil { + t.Fatal("expected error when extra_address is missing") + } + }) + + t.Run("extra address without listener", func(t *testing.T) { + cfg := validSyslog(t) + cfg.ExtraAddress = "127.0.0.1:5144" + if err := cfg.Validate(); err == nil { + t.Fatal("expected error when extra_listener is missing") + } + }) + + t.Run("duplicate listener", func(t *testing.T) { + cfg := validSyslog(t) + cfg.Listener = TCPListener + cfg.Address = "127.0.0.1:5144" + cfg.ExtraListener = TCPListener + cfg.ExtraAddress = "127.0.0.1:5144" + if err := cfg.Validate(); err == nil { + t.Fatal("expected error when extra listener duplicates the primary") + } + }) + + t.Run("invalid extra listener", func(t *testing.T) { + cfg := validSyslog(t) + cfg.ExtraListener = "sctp" + cfg.ExtraAddress = "127.0.0.1:5144" + if err := cfg.Validate(); err == nil { + t.Fatal("expected error for invalid extra_listener") + } + }) +} diff --git a/syslog/syslog.go b/syslog/syslog.go index 75ba3da..c32fc1a 100644 --- a/syslog/syslog.go +++ b/syslog/syslog.go @@ -106,45 +106,69 @@ func (s *SyslogWorker) Start() error { return errors.Wrap(err, "removing socket") } - switch s.cfg.Listener { + if err := s.listen(s.cfg.Listener, s.cfg.Address); err != nil { + return err + } + if s.cfg.ExtraListener != "" { + if err := s.listen(s.cfg.ExtraListener, s.cfg.ExtraAddress); err != nil { + return err + } + } + + err := s.server.Boot() + if err != nil { + return errors.Wrap(err, "starting syslog server") + } + go s.doWork() + return nil +} + +func (s *SyslogWorker) listen(listener config.ListenerType, address string) error { + switch listener { case config.UnixDgramListener: - if err := s.server.ListenUnixgram(s.cfg.Address); err != nil { - return errors.Wrap(err, fmt.Sprintf("listening on unix socket %q", s.cfg.Address)) + if err := s.server.ListenUnixgram(address); err != nil { + return errors.Wrap(err, fmt.Sprintf("listening on unix socket %q", address)) } - if _, err := os.Stat(s.cfg.Address); err != nil { - log.Warningf("cannot fetch info about %q: %q", s.cfg.Address, err) + if _, err := os.Stat(address); err != nil { + log.Warningf("cannot fetch info about %q: %q", address, err) } else { - if err := os.Chmod(s.cfg.Address, 0666); err != nil { - log.Warningf("cannot change permissions on %q: %q", s.cfg.Address, err) + if err := os.Chmod(address, 0666); err != nil { + log.Warningf("cannot change permissions on %q: %q", address, err) } } case config.TCPListener: - if err := s.server.ListenTCP(s.cfg.Address); err != nil { - return errors.Wrap(err, fmt.Sprintf("listening on TCP %q", s.cfg.Address)) + if err := s.server.ListenTCP(address); err != nil { + return errors.Wrap(err, fmt.Sprintf("listening on TCP %q", address)) } case config.UDPListener: - if err := s.server.ListenUDP(s.cfg.Address); err != nil { - return errors.Wrap(err, fmt.Sprintf("listening on UDP %q", s.cfg.Address)) + if err := s.server.ListenUDP(address); err != nil { + return errors.Wrap(err, fmt.Sprintf("listening on UDP %q", address)) } + default: + return fmt.Errorf("invalid listener type %q", listener) } + return nil +} - err := s.server.Boot() - if err != nil { - return errors.Wrap(err, "starting syslog server") +func (s *SyslogWorker) unixAddresses() []string { + var addrs []string + if s.cfg.Listener == config.UnixDgramListener { + addrs = append(addrs, s.cfg.Address) } - go s.doWork() - return nil + if s.cfg.ExtraListener == config.UnixDgramListener { + addrs = append(addrs, s.cfg.ExtraAddress) + } + return addrs } func (s *SyslogWorker) cleanStaleSocket() error { - if s.cfg.Listener != config.UnixDgramListener { - return nil - } - if mode, err := os.Stat(s.cfg.Address); err == nil { - if mode.Mode()&os.ModeSocket != 0 { - log.Infof("removing unix socket %q", s.cfg.Address) - if err := os.Remove(s.cfg.Address); err != nil { - return errors.Wrap(err, "removing unix socket") + for _, addr := range s.unixAddresses() { + if mode, err := os.Stat(addr); err == nil { + if mode.Mode()&os.ModeSocket != 0 { + log.Infof("removing unix socket %q", addr) + if err := os.Remove(addr); err != nil { + return errors.Wrap(err, "removing unix socket") + } } } } diff --git a/syslog/syslog_test.go b/syslog/syslog_test.go new file mode 100644 index 0000000..cc5be97 --- /dev/null +++ b/syslog/syslog_test.go @@ -0,0 +1,157 @@ +// Copyright 2026 Cloudbase Solutions SRL +// +// Licensed under the Apache License, Version 2.0 (the "License"); you may +// not use this file except in compliance with the License. You may obtain +// a copy of the License at +// +// http://www.apache.org/licenses/LICENSE-2.0 +// +// Unless required by applicable law or agreed to in writing, software +// distributed under the License is distributed on an "AS IS" BASIS, WITHOUT +// WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. See the +// License for the specific language governing permissions and limitations +// under the License. + +package syslog + +import ( + "context" + "fmt" + "net" + "path/filepath" + "sync" + "testing" + "time" + + "coriolis-logger/config" + "coriolis-logger/logging" +) + +type collectingWriter struct { + mu sync.Mutex + msgs []logging.LogMessage +} + +func (c *collectingWriter) Write(logMsg logging.LogMessage) error { + c.mu.Lock() + defer c.mu.Unlock() + c.msgs = append(c.msgs, logMsg) + return nil +} + +func (c *collectingWriter) messages() []logging.LogMessage { + c.mu.Lock() + defer c.mu.Unlock() + out := make([]logging.LogMessage, len(c.msgs)) + copy(out, c.msgs) + return out +} + +func freeListenAddr(t *testing.T) string { + t.Helper() + ln, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("reserving tcp port: %v", err) + } + addr := ln.Addr().String() + if err := ln.Close(); err != nil { + t.Fatalf("releasing tcp port: %v", err) + } + return addr +} + +func syslogRFC5424(msg string) string { + return fmt.Sprintf("<14>1 2024-01-02T03:04:05Z testhost coriolis-logger - - - %s", msg) +} + +func sendUnixgram(t *testing.T, socket, msg string) { + t.Helper() + conn, err := net.Dial("unixgram", socket) + if err != nil { + t.Fatalf("dial unixgram: %v", err) + } + defer conn.Close() + if _, err := conn.Write([]byte(msg)); err != nil { + t.Fatalf("write unixgram: %v", err) + } +} + +func sendTCP(t *testing.T, addr, msg string) { + t.Helper() + var conn net.Conn + var err error + deadline := time.Now().Add(2 * time.Second) + for { + conn, err = net.DialTimeout("tcp", addr, 200*time.Millisecond) + if err == nil { + break + } + if time.Now().After(deadline) { + t.Fatalf("dial tcp: %v", err) + } + time.Sleep(20 * time.Millisecond) + } + defer conn.Close() + if _, err := conn.Write([]byte(msg + "\n")); err != nil { + t.Fatalf("write tcp: %v", err) + } +} + +func waitForMessages(t *testing.T, writer *collectingWriter, n int) []logging.LogMessage { + t.Helper() + deadline := time.Now().Add(3 * time.Second) + for time.Now().Before(deadline) { + msgs := writer.messages() + if len(msgs) >= n { + return msgs + } + time.Sleep(20 * time.Millisecond) + } + t.Fatalf("timed out waiting for %d messages, got %d", n, len(writer.messages())) + return nil +} + +func TestSyslogWorkerUnixAndTCP(t *testing.T) { + socket := filepath.Join(t.TempDir(), "syslog.sock") + tcpAddr := freeListenAddr(t) + + cfg := config.Syslog{ + Listener: config.UnixDgramListener, + Address: socket, + ExtraListener: config.TCPListener, + ExtraAddress: tcpAddr, + Format: "rfc5424", + DataStore: config.StdOutDataStore, + } + + ctx, cancel := context.WithCancel(context.Background()) + defer cancel() + + writer := &collectingWriter{} + worker, err := NewSyslogServer(ctx, cfg, writer, make(chan error, 1)) + if err != nil { + t.Fatalf("creating syslog server: %v", err) + } + if err := worker.Start(); err != nil { + t.Fatalf("starting syslog server: %v", err) + } + defer func() { + cancel() + worker.Wait() + }() + + sendUnixgram(t, socket, syslogRFC5424("from-unix")) + sendTCP(t, tcpAddr, syslogRFC5424("from-tcp")) + + msgs := waitForMessages(t, writer, 2) + seen := map[string]bool{} + for _, msg := range msgs { + seen[msg.Message] = true + } + if !seen["from-unix"] { + t.Errorf("missing unix datagram message, got %+v", msgs) + } + if !seen["from-tcp"] { + t.Errorf("missing tcp message, got %+v", msgs) + } +} diff --git a/testdata/config.toml b/testdata/config.toml index b72f9ec..ee03bba 100644 --- a/testdata/config.toml +++ b/testdata/config.toml @@ -36,6 +36,12 @@ listener = "unixgram" # address = "/tmp/coriolis-logger/syslog" address = "/tmp/coriolis-logging.sock" +# Optional second listener. The syslog server can bind a unix +# datagram socket and a TCP or UDP endpoint at the same time. +# extra_listener and extra_address must both be set when used. +# extra_listener = "tcp" +# extra_address = "0.0.0.0:5144" + # Log format # possible values: # rfc3164