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
2 changes: 1 addition & 1 deletion api/http.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ func (h *Handler) log(w http.ResponseWriter, req *http.Request) {
return
}
logger := log.WithFunc("api.log").WithField("path", "/log")
// the status line must go out before the hijack, otherwise clients see no response
// The status line must go out before the hijack, otherwise clients see no response
w.WriteHeader(http.StatusOK)
if hijack, ok := w.(http.Hijacker); ok {
conn, buf, err := hijack.Hijack()
Expand Down
2 changes: 1 addition & 1 deletion collector/collector.go
Original file line number Diff line number Diff line change
Expand Up @@ -81,7 +81,7 @@ func New(ctx context.Context, config *types.Config) *Collector {

total, err := hostMemTotal(c.procRoot)
if err != nil {
// without a node total the memory percentages of an unlimited workload stay unreported
// Without a node total the memory percentages of an unlimited workload stay unreported
log.WithFunc("collector.New").Warnf(ctx, "failed to read the node memory total: %v", err)
}
c.memTotal = total
Expand Down
2 changes: 1 addition & 1 deletion collector/console.go
Original file line number Diff line number Diff line change
Expand Up @@ -66,7 +66,7 @@ func (c *Console) Read(ctx context.Context, handle EntryHandler) {
if delivered {
backoff = consoleRetryMin
}
// the console goes away with the vm and comes back with it, so a closed one is not a failure
// The console goes away with the vm and comes back with it, so a closed one is not a failure
logger.Debugf(ctx, "console stopped: %v", err)
if c.dropped > 0 {
logger.Warnf(ctx, "the journal refused %d console lines", c.dropped)
Expand Down
8 changes: 4 additions & 4 deletions collector/health.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ package collector

import (
"context"
"fmt"
"net"
"sync"
"time"

Expand All @@ -16,18 +16,18 @@ func Probe(ctx context.Context, w *source.Workload, timeout time.Duration) bool
if check == nil {
return true
}
// without an address of its own there is nothing to dial, and the node's own ports are not the workload's
// Without an address of its own there is nothing to dial, and the node's own ports are not the workload's
if w.LocalIP == "" {
return false
}

var tcpChecker []string
for _, port := range check.TCPPorts {
tcpChecker = append(tcpChecker, fmt.Sprintf("%s:%s", w.LocalIP, port))
tcpChecker = append(tcpChecker, net.JoinHostPort(w.LocalIP, port))
}
httpURL := ""
if check.HTTPPort != "" {
httpURL = fmt.Sprintf("http://%s:%s%s", w.LocalIP, check.HTTPPort, check.HTTPURL)
httpURL = "http://" + net.JoinHostPort(w.LocalIP, check.HTTPPort) + check.HTTPURL
}

var httpOK, tcpOK bool
Expand Down
24 changes: 24 additions & 0 deletions collector/health_test.go
Original file line number Diff line number Diff line change
Expand Up @@ -2,6 +2,8 @@ package collector

import (
"net"
"net/http"
"net/http/httptest"
"testing"
"time"

Expand Down Expand Up @@ -41,6 +43,28 @@ func TestProbeAcceptsAnOpenTCPPort(t *testing.T) {
assert.True(t, Probe(t.Context(), w, time.Second))
}

func TestProbeReachesAnIPv6Workload(t *testing.T) {
listener, err := net.Listen("tcp", "[::1]:0")
if err != nil {
t.Skipf("no ipv6 loopback: %v", err)
}
server := httptest.NewUnstartedServer(http.HandlerFunc(func(http.ResponseWriter, *http.Request) {}))
_ = server.Listener.Close()
server.Listener = listener
server.Start()
defer server.Close()

_, port, err := net.SplitHostPort(listener.Addr().String())
require.NoError(t, err)

w := &source.Workload{
ID: "ipv6",
LocalIP: "::1",
Meta: source.Meta{HealthCheck: &coretypes.HealthCheck{TCPPorts: []string{port}, HTTPPort: port, HTTPURL: "/"}},
}
assert.True(t, Probe(t.Context(), w, time.Second))
}

func TestProbeRejectsAClosedTCPPort(t *testing.T) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
Expand Down
4 changes: 2 additions & 2 deletions collector/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -109,7 +109,7 @@ func (j *Journal) Read(ctx context.Context, handle EntryHandler) error {
logger.Error(ctx, err, "failed to decode a journal record")
continue
}
// a console line reached the journal through the reader that already forwarded it
// A console line reached the journal through the reader that already forwarded it
if record.EruStream != common.StreamConsole {
handle(record.entry())
}
Expand All @@ -125,7 +125,7 @@ func (j *Journal) Read(ctx context.Context, handle EntryHandler) error {
_ = cmd.Wait()
return err
}
// a journalctl that dies looks like a clean eof on stdout, so its status is the only signal
// A journalctl that dies looks like a clean eof on stdout, so its status is the only signal
if err := cmd.Wait(); err != nil {
return fmt.Errorf("%s exited: %w: %s", j.binary, err, stderr)
}
Expand Down
3 changes: 0 additions & 3 deletions logs/writer.go
Original file line number Diff line number Diff line change
Expand Up @@ -19,7 +19,6 @@ const (
Discard = "__discard__"

keepaliveInterval = time.Second * 30
closeWaitInterval = time.Second * 5
dialTimeout = time.Second * 5
writeTimeout = time.Second * 5
)
Expand Down Expand Up @@ -157,8 +156,6 @@ func (w *Writer) keepalive(ctx context.Context) {
case <-tick.C:
w.reconnect(ctx)
case <-ctx.Done():
// give the pending writes a chance to drain before closing
time.Sleep(closeWaitInterval)
if err := w.close(ctx); err != nil {
log.WithFunc("logs.keepalive").Errorf(ctx, err, "failed to close writer %s", w.addr)
}
Expand Down
25 changes: 22 additions & 3 deletions logs/writer_test.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,12 @@
package logs

import (
"context"
"errors"
"io"
"net"
"strings"
"sync"
"testing"
"time"

Expand Down Expand Up @@ -76,8 +78,9 @@ func TestNewWriters(t *testing.T) {

ctx := t.Context()

var writers sync.WaitGroup
for addr, expectedErr := range cases {
go func(addr string, expectedErr error) {
writers.Go(func() {
writer, err := NewWriter(ctx, addr, false)
assert.Equal(t, expectedErr, err)
if expectedErr != nil {
Expand All @@ -86,9 +89,25 @@ func TestNewWriters(t *testing.T) {
assert.NoError(t, err)
err = writer.Write(ctx, &types.Log{})
assert.NoError(t, err)
}(addr, expectedErr)
})
}
time.Sleep(closeWaitInterval + 2*time.Second)
writers.Wait()
}

func TestWriterClosesOnceItsContextIsCancelled(t *testing.T) {
listener, err := net.Listen("tcp", "127.0.0.1:0")
require.NoError(t, err)
defer listener.Close()

ctx, cancel := context.WithCancel(t.Context())
w, err := NewWriter(ctx, "tcp://"+listener.Addr().String(), false)
require.NoError(t, err)
require.NoError(t, w.Write(ctx, &types.Log{}))

cancel()
assert.Eventually(t, func() bool {
return errors.Is(w.Write(ctx, &types.Log{}), common.ErrConnecting)
}, time.Second, 10*time.Millisecond)
}

func TestWriteKeepsTheEncoderWhenADatagramIsTooBig(t *testing.T) {
Expand Down
2 changes: 1 addition & 1 deletion logshim/logshim.go
Original file line number Diff line number Diff line change
Expand Up @@ -105,7 +105,7 @@ func (s *stream) pump(reader io.Reader) {
if err != nil {
return
}
// a journal that refuses a line must not fail the container, so the line is lost
// A journal that refuses a line must not fail the container, so the line is lost
if sendErr := s.send(string(line), s.priority, vars); sendErr != nil {
s.dropped++
s.err = sendErr
Expand Down
2 changes: 1 addition & 1 deletion manager/node/heartbeat.go
Original file line number Diff line number Diff line change
Expand Up @@ -40,7 +40,7 @@ func (m *Manager) nodeStatusReport(ctx context.Context) {
return
}

// the ttl outlives the interval so one lost report cannot expire the node
// The ttl outlives the interval so one lost report cannot expire the node
ttl := int64(m.config.HeartbeatInterval * ttlHeartbeats)

if err := utils.BackoffRetry(ctx, reportAttempts, func() error {
Expand Down
2 changes: 1 addition & 1 deletion manager/node/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -39,7 +39,7 @@ func (m *Manager) Exit(ctx context.Context) error {
logger := log.WithFunc("node.Exit").WithField("hostname", m.config.HostName)
logger.Info(ctx, "remove node status")

// a negative ttl removes the node status
// A negative ttl removes the node status
return m.setNodeStatus(ctx, -1)
}

Expand Down
2 changes: 1 addition & 1 deletion manager/workload/event.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,7 @@ func (e *EventHandler) Watch(ctx context.Context, c <-chan *types.WorkloadEventM
return
}
logger.Infof(ctx, "workload %s action %s", ev.ID, ev.Action)
// one workload's events are applied in order, so a die cannot land after the next start
// One workload's events are applied in order, so a die cannot land after the next start
switch ev.Action {
case common.StatusStart:
e.queue.Go(ev.ID, func() { e.start(ctx, ev) })
Expand Down
2 changes: 1 addition & 1 deletion manager/workload/journal.go
Original file line number Diff line number Diff line change
Expand Up @@ -47,7 +47,7 @@ func (m *Manager) forwardJournal(ctx context.Context) {

// forwardConsole reads a vm's serial console, which journald only holds once this reader has written it there.
func (m *Manager) forwardConsole(ctx context.Context, w *source.Workload) {
// a restarted vm gets a new console, so the path is read back per attempt rather than held from here
// A restarted vm gets a new console, so the path is read back per attempt rather than held from here
console := collector.NewConsole(w.ID, w.Meta.Appname, func() (string, error) {
fresh := m.refreshed(ctx, w.ID)
if fresh == nil {
Expand Down
18 changes: 7 additions & 11 deletions manager/workload/manager.go
Original file line number Diff line number Diff line change
Expand Up @@ -76,14 +76,14 @@ func NewManager(ctx context.Context, config *types.Config, clients *manager.Clie
}

func (m *Manager) Run(ctx context.Context) error {
// watching before the initial load means an event raised during it is handled, not missed
// Watching before the initial load means an event raised during it is handled, not missed
go m.monitor(ctx)

if err := m.initWorkloadStatus(ctx); err != nil {
return err
}

// the journal reader starts once the load registered every target, so no backlog line is dropped
// The journal reader starts once the load registered every target, so no backlog line is dropped
go m.forwardJournal(ctx)

go m.healthCheck(ctx)
Expand All @@ -97,15 +97,11 @@ func (m *Manager) PullLog(ctx context.Context, app string, buf *bufio.ReadWriter
ID, errChan, unsubscribe := m.logBroadcaster.subscribe(ctx, app, buf)
defer unsubscribe()

for {
select {
case <-ctx.Done():
return
case err := <-errChan:
if !errors.Is(err, io.EOF) {
log.WithFunc("workload.PullLog").WithField("ID", ID).Error(ctx, err, "failed to pull log")
}
return
select {
case <-ctx.Done():
case err := <-errChan:
if !errors.Is(err, io.EOF) {
log.WithFunc("workload.PullLog").WithField("ID", ID).Error(ctx, err, "failed to pull log")
}
}
}
Expand Down
2 changes: 1 addition & 1 deletion ocihook/ocihook.go
Original file line number Diff line number Diff line change
Expand Up @@ -67,7 +67,7 @@ func Command() *cli.Command {
},
},
Action: func(ctx context.Context, cmd *cli.Command) error {
// the hook runs inside runc's create, which waits for it however long an ipam plugin takes
// The hook runs inside runc's create, which waits for it however long an ipam plugin takes
ctx, cancel := context.WithTimeout(ctx, timeout)
defer cancel()

Expand Down
2 changes: 1 addition & 1 deletion source/cocoon/cocoon.go
Original file line number Diff line number Diff line change
Expand Up @@ -107,7 +107,7 @@ func (c *Cocoon) watchDaemon(ctx context.Context) {
logger := log.WithFunc("cocoon.watchDaemon")
for {
err := c.daemon.events(ctx, func(ID string, running, gone bool) {
// the daemon reports every vm on the node, not only eru's
// The daemon reports every vm on the node, not only eru's
if !meta.IsID(ID) {
return
}
Expand Down
4 changes: 2 additions & 2 deletions source/containerd/events.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,15 +78,15 @@ func translate(ctx context.Context, envelope *events.Envelope) *types.WorkloadEv
case *apievents.TaskStart:
ID, action = e.ContainerID, common.StatusStart
case *apievents.TaskExit:
// an exec process exiting is not the workload exiting
// An exec process exiting is not the workload exiting
if e.ID != e.ContainerID {
return nil
}
ID, action = e.ContainerID, common.StatusDie
case *apievents.ContainerDelete:
ID, action = e.ID, common.StatusDie
case *apievents.ContainerUpdate:
// the oci hook writes the cni addresses back as labels, so an update is a new set of facts
// The oci hook writes the cni addresses back as labels, so an update is a new set of facts
ID, action = e.ID, common.StatusStart
default:
return nil
Expand Down
4 changes: 2 additions & 2 deletions source/containerd/meta.go
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@ func readSpec(raw typeurl.Any) (spec, error) {
if oci.Process != nil {
s.env = oci.Process.Env
}
// a container that shares the node's network has no network namespace of its own in the spec
// A container that shares the node's network has no network namespace of its own in the spec
s.hostNetwork = oci.Linux != nil && !slices.ContainsFunc(oci.Linux.Namespaces, func(ns ociNamespace) bool {
return ns.Type == string(specs.NetworkNamespace)
})
Expand All @@ -73,7 +73,7 @@ func (c *Containerd) workload(ctx context.Context, ID string, labels map[string]
meta := coreutils.DecodeMetaInLabel(ctx, labels)
nets := networks(labels)

// core's engines all report the node's own address for a host network workload
// Core's engines all report the node's own address for a host network workload
localIP := source.Addr(nets)
if len(nets) == 0 && s.hostNetwork {
nets, localIP = map[string]string{hostNetwork: c.nodeIP}, common.LocalIP
Expand Down
1 change: 0 additions & 1 deletion source/meta/meta.go
Original file line number Diff line number Diff line change
Expand Up @@ -26,7 +26,6 @@ const (
var (
errOtherKind = errors.New("meta file of another runtime")

// workloadID is the shape of a workload id, so nothing else eru names on a node is taken for one.
workloadID = regexp.MustCompile("^[0-9a-f]{32}$")
)

Expand Down
2 changes: 1 addition & 1 deletion source/multi.go
Original file line number Diff line number Diff line change
Expand Up @@ -78,7 +78,7 @@ func (m *Multi) Events(ctx context.Context) (<-chan *types.WorkloadEventMessage,
defer close(eventChan)
defer close(errChan)

// one runtime failing stops them all, so the manager resubscribes to every runtime at once
// One runtime failing stops them all, so the manager resubscribes to every runtime at once
var wg sync.WaitGroup
for _, src := range m.sources {
events, errs := src.Events(ctx)
Expand Down
8 changes: 4 additions & 4 deletions source/systemd/systemd.go
Original file line number Diff line number Diff line change
Expand Up @@ -116,7 +116,7 @@ func (s *Systemd) Alive(ctx context.Context) bool {

func (s *Systemd) watchUnits(ctx context.Context) error {
logger := log.WithFunc("systemd.watchUnits")
// a previous Events left its subscription on this shared connection, still feeding dead channels
// A previous Events left its subscription on this shared connection, still feeding dead channels
_ = s.conn.Unsubscribe()
if err := s.conn.Subscribe(); err != nil {
return err
Expand All @@ -133,7 +133,7 @@ func (s *Systemd) watchUnits(ctx context.Context) error {
case update := <-updates:
ID, ok := workloadIDFromUnit(update.UnitName)
if !ok {
// the subscription is node wide, so only a name under eru's prefix is worth a line
// The subscription is node wide, so only a name under eru's prefix is worth a line
if strings.HasPrefix(update.UnitName, unitPrefix) {
logger.Debugf(ctx, "ignoring unit %s, it is not a workload", update.UnitName)
}
Expand All @@ -147,7 +147,7 @@ func (s *Systemd) watchUnits(ctx context.Context) error {
s.reporter.Report(ID, action)
}
case err := <-errs:
// a subscriber that fell behind missed transitions, it did not lose the bus
// A subscriber that fell behind missed transitions, it did not lose the bus
logger.Warnf(ctx, "systemd subscription fell behind, relisting: %v", err)
if err := s.relist(ctx); err != nil {
return err
Expand All @@ -171,7 +171,7 @@ func (s *Systemd) relist(ctx context.Context) error {
}

func (s *Systemd) runningUnits(ctx context.Context) (map[string]bool, error) {
// the bus matches a glob only, so eru-agent.service comes back too and is dropped here
// The bus matches a glob only, so eru-agent.service comes back too and is dropped here
units, err := s.conn.ListUnitsByPatternsContext(ctx, nil, []string{unitPattern})
if err != nil {
return nil, err
Expand Down
1 change: 0 additions & 1 deletion store/core/node.go
Original file line number Diff line number Diff line change
Expand Up @@ -18,7 +18,6 @@ func (s *Store) GetNode(ctx context.Context, nodename string) (*types.Node, erro
return &types.Node{Endpoint: resp.Endpoint}, nil
}

// SetNodeStatus reports the node alive under core's ttl.
func (s *Store) SetNodeStatus(ctx context.Context, ttl int64) error {
opts := &pb.SetNodeStatusOptions{
Nodename: s.config.HostName,
Expand Down
3 changes: 1 addition & 2 deletions store/core/workload.go
Original file line number Diff line number Diff line change
Expand Up @@ -33,7 +33,6 @@ func (s *Store) ListRunningWorkloadIDs(ctx context.Context) ([]string, error) {
return IDs, nil
}

// WorkloadExists reports whether core still owns the workload.
func (s *Store) WorkloadExists(ctx context.Context, ID string) (bool, error) {
_, err := call(ctx, s, func(ctx context.Context) (*pb.Workload, error) {
return s.client().GetWorkload(ctx, &pb.WorkloadID{Id: ID})
Expand All @@ -47,7 +46,7 @@ func (s *Store) SetWorkloadStatus(ctx context.Context, status *types.WorkloadSta
return nil
}

// core's selfmon owns status expiry
// Core's selfmon owns status expiry
statusPb := &pb.WorkloadStatus{
Id: status.ID,
Running: status.Running,
Expand Down
Loading
Loading