Skip to content

Commit 032865d

Browse files
authored
Count dropped log lines per point (#126)
* test: benchmark the forward writer per sink Discard, loopback tcp and loopback udp, so the serial forward path has a measured per-line ceiling to compare future changes against. * feat: count dropped log lines per point log_lines_dropped_total{point=forward|subscriber|console} on /metrics. Every drop on the agent side was already counted and warned about; none of it was observable without reading logs. The shim stays out: it is a separate process and reports its drops once when it exits. * feat: count what journald rate-limits journald drops a service's excess lines inside itself and tells the sender nothing, so the shim and console counters never see it; the reader ORs journald's SD_MESSAGE_JOURNAL_DROPPED notices into its match and books N_DROPPED for eru units and the agent's own senders as point=journald. On a stock node a 60000-line burst keeps 37500. The counter now registers where it is declared, so a node without an api addr counts the same way. * review: apply the simplify round forward keeps its one-line error shape, the test helper carries a broadcaster, the counter asserts compare exactly, and the writer benchmark swaps the no-op arm for an encode arm so the agent's own per-line cost is visible next to the socket floor. * perf: encode forward lines into a reused buffer json.Marshal returns cap==len, so appending the newline reallocated and copied every line. A buffer and json.Encoder kept on the StreamEncoder cut the encode from 480 B/5 allocs to 64 B/4 allocs per line; the tcp arm of BenchmarkWriterWrite went from 6.5-7.7 us to 4.3-4.6 us. * fix: stop inferring journald log drops Journald suppression notices are service-wide, delayed until another allowed message, and cannot close a per-workload counter. Keep the metric limited to drops the agent observes directly. * docs: keep the journald rate limit and the no-blocking property The revert of the journald point also removed the operator facts that stand on their own: journald drops a fast service's lines silently at its own per-service limit, a unit's LogRateLimitBurst= is the override, and nothing on the log path blocks a workload. * docs: narrow the log drop guarantee Journald writes have no deadline, so the README must not promise that every logging failure is non-blocking. Keep the claim limited to observable drops. * review: say what the shim guards against The comment promised the container never blocks on the journal; the code only drops a line the journal refuses, and a journal that stops draining its socket still back-pressures the container (journal.Send has no deadline).
1 parent 32fdc24 commit 032865d

14 files changed

Lines changed: 153 additions & 16 deletions

File tree

‎README.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,9 +11,9 @@ Eru's per-node agent. It watches the workloads a node runs — containerd contai
1111

1212
- **Node heartbeat** — reports the node alive to core on an interval, with a ttl three times the interval so one lost report does not evict the node; a clean shutdown removes the status, a `SIGUSR1` restart keeps it.
1313
- **Workload health checks** — TCP and HTTP probes declared per workload, run on every runtime event and on a periodic sweep, with the result published through core.
14-
- **Log forwarding** — ships each line as a JSON record to `tcp://`, `udp://` or `journal://` targets, sharded over several targets by workload id, read from the node's journal with a persisted cursor, or read straight off a VM's serial console. Container output reaches the journal through `eru-agent log-shim`, containerd's binary logger; console output is written there by the agent itself, so every workload kind has the same history.
14+
- **Log forwarding** — ships each line as a JSON record to `tcp://`, `udp://` or `journal://` targets, sharded over several targets by workload id, read from the node's journal with a persisted cursor, or read straight off a VM's serial console. Container output reaches the journal through `eru-agent log-shim`, containerd's binary logger; console output is written there by the agent itself, so every workload kind has the same history. A down forward target, a tailing client that falls behind, or journald's own per-service rate limit can drop lines; the agent counts only the drops it can observe.
1515
- **Live log tailing** — `GET /log/?app=<name>` streams the logs of one application straight off the node.
16-
- **Prometheus metrics** — per-workload cpu, memory, per-nic network and per-device block io gauges on `/metrics`, optionally pushed to statsd as well. Sampled straight from cgroup v2 files, so a tick makes no call to any daemon.
16+
- **Prometheus metrics** — per-workload cpu, memory, per-nic network and per-device block io gauges on `/metrics`, optionally pushed to statsd as well, plus a node-level `log_lines_dropped_total{point}` counter for the log path. Sampled straight from cgroup v2 files, so a tick makes no call to any daemon.
1717
- **Runtimes per node** — a node declares the runtimes it hosts under a required `runtimes:` section: containerd containers, systemd process pods, cocoon VMs, or several at once; the heartbeat needs every one of them alive. The old `runtime:` and `docker:` keys are gone, with no compatibility shim. A mock runtime and store cover development.
1818
- **Node-side helper modes** — the same binary is containerd's task logger (`eru-agent log-shim`) and its CNI hook (`eru-agent oci-hook`), so a container node needs no eru daemon beyond the agent.
1919

‎collector/console.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ const (
2222
consoleRetryMax = 5 * time.Second
2323
)
2424

25+
var droppedByConsole = LogLinesDropped.WithLabelValues(DropPointConsole)
26+
2527
type journalFunc func(message string, priority journal.Priority, vars map[string]string) error
2628

2729
// pathFunc answers where the vm's console is now: a restart moves it, and core rewrites the meta file.
@@ -129,5 +131,6 @@ func (c *Console) emit(line string, handle EntryHandler) {
129131
// journald holds the history core reads back over ssh, the same way it does for a container
130132
if err := c.send(line, journal.PriInfo, c.vars); err != nil {
131133
c.dropped++
134+
droppedByConsole.Inc()
132135
}
133136
}

‎collector/console_test.go‎

Lines changed: 3 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -12,6 +12,7 @@ import (
1212
"time"
1313

1414
"github.com/coreos/go-systemd/v22/journal"
15+
"github.com/prometheus/client_golang/prometheus/testutil"
1516
"github.com/stretchr/testify/assert"
1617
"github.com/stretchr/testify/require"
1718

@@ -154,13 +155,15 @@ func TestEmitStampsTheWorkloadAndTheConsoleStream(t *testing.T) {
154155
func TestEmitKeepsForwardingWhenTheJournalRefuses(t *testing.T) {
155156
console := NewConsole(consoleWorkload, consoleAppname, at("unused"))
156157
console.send = func(string, journal.Priority, map[string]string) error { return errJournalRefused }
158+
counted := testutil.ToFloat64(droppedByConsole)
157159

158160
var entries []*Entry
159161
console.emit("boot", func(e *Entry) { entries = append(entries, e) })
160162
console.emit("login:", func(e *Entry) { entries = append(entries, e) })
161163

162164
assert.Len(t, entries, 2)
163165
assert.Equal(t, 2, console.dropped)
166+
assert.Equal(t, counted+2, testutil.ToFloat64(droppedByConsole))
164167
}
165168

166169
func at(path string) pathFunc {

‎collector/metrics.go‎

Lines changed: 13 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -11,18 +11,29 @@ import (
1111
"github.com/projecteru2/core/log"
1212
coreutils "github.com/projecteru2/core/utils"
1313
"github.com/prometheus/client_golang/prometheus"
14+
"github.com/prometheus/client_golang/prometheus/promauto"
1415

1516
"github.com/projecteru2/agent/source"
1617
)
1718

1819
const (
19-
labelNIC = "nic"
20-
labelDev = "dev"
20+
DropPointSubscriber = "subscriber"
21+
DropPointConsole = "console"
22+
DropPointForward = "forward"
23+
24+
labelNIC = "nic"
25+
labelDev = "dev"
26+
labelPoint = "point"
2127

2228
metricMemMaxUsage = "mem_max_usage"
2329
)
2430

2531
var (
32+
LogLinesDropped = promauto.With(prometheus.DefaultRegisterer).NewCounterVec(prometheus.CounterOpts{
33+
Name: "log_lines_dropped_total",
34+
Help: "workload log lines that never reached their destination, by the point that lost them.",
35+
}, []string{labelPoint})
36+
2637
clientsMutex sync.Mutex
2738
clients = map[string]*MetricsClient{}
2839

‎docs/architecture.md‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -80,7 +80,7 @@ Either way, each line becomes the same JSON record:
8080

8181
The record goes to two places: the configured forwarder for that workload, and the in-process broadcaster that serves `/log/`. Non-utf8 bytes are escaped as `\xNN` so a binary blob on stdout cannot corrupt the stream.
8282

83-
**Log broadcaster.** `GET /log/?app=<name>` hijacks the connection and subscribes to every record whose `name` matches. Each reader broadcasts its lines in order, and one workload's lines come from one reader, so a subscriber sees a workload's lines in order. A subscriber whose connection breaks — detected by a read on the hijacked socket — is cancelled and dropped.
83+
**Log broadcaster.** `GET /log/?app=<name>` hijacks the connection and subscribes to every record whose `name` matches. Each reader broadcasts its lines in order, and one workload's lines come from one reader, so a subscriber sees a workload's lines in order. A subscriber whose connection breaks — detected by a read on the hijacked socket — is cancelled and dropped; one that falls more than 256 lines behind loses the excess rather than stalling the node's forwarding, counted in `log_lines_dropped_total{point="subscriber"}`.
8484

8585
## Status reporting
8686

‎docs/configuration.md‎

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -90,10 +90,12 @@ log:
9090
stdout: false
9191
```
9292

93-
`forwards` is a list of targets. Supported schemes are `tcp://`, `udp://` and `journal://`; a target with any other scheme is accepted and silently discards, with a warning at startup. Each workload is pinned to one target by hashing its id, so several targets share the load without duplicating lines. A target that is down is retried every 30 seconds in the background while its lines are dropped; a line too large for a udp datagram loses only itself.
93+
`forwards` is a list of targets. Supported schemes are `tcp://`, `udp://` and `journal://`; a target with any other scheme is accepted and silently discards, with a warning at startup. Each workload is pinned to one target by hashing its id, so several targets share the load without duplicating lines. A target that is down is retried every 30 seconds in the background while its lines are dropped; a line too large for a udp datagram loses only itself. Both count in `log_lines_dropped_total{point="forward"}` on `/metrics`.
9494

9595
A VM's output is its serial console, so the agent reads the console its meta file names, one goroutine per VM, forwarding each line and writing it to journald so the history is there too. Every other runtime logs to journald natively, and the agent runs one `journalctl --follow --output=json SYSLOG_IDENTIFIER=eru` for the whole node, resuming from the cursor it saved under `state_dir`. That path needs `journalctl` on the node, and needs every eru workload to log under the `eru` syslog identifier: process units get it from core's `systemd-run`, containers from `eru-agent log-shim`.
9696

97+
journald rate-limits every service on its own — by default `RateLimitBurst=10000` per `RateLimitIntervalSec=30s`, per priority class, scaled up with the journal's free space — and drops the excess without telling the sender, so a workload that logs faster than that loses lines before the agent ever sees them. A unit's own `LogRateLimitBurst=` overrides the node default; journald's "Suppressed N messages" notices are the only record of what was lost.
98+
9799
`stdout: true` additionally writes every forwarded line to the agent's own log.
98100

99101
## `metrics`

‎docs/metrics.md‎

Lines changed: 18 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -16,6 +16,24 @@ Sampling happens every `metrics.step` seconds. Rates are computed over that wind
1616

1717
The node-wide cpu split that the two `cpu_host_*_usage` ratios divide by is read from `/proc/stat` at most once a second and shared by every workload, instead of once per workload per tick.
1818

19+
## Counters
20+
21+
One node-level counter, with no workload labels:
22+
23+
| Counter | Meaning |
24+
|---|---|
25+
| `log_lines_dropped_total{point}` | workload log lines the agent observed losing, by where they were lost |
26+
27+
| `point` | when |
28+
|---|---|
29+
| `forward` | the workload's forward target was down, timed out on a write, or refused an oversized datagram; the line is not replayed after the reconnect |
30+
| `subscriber` | a `GET /log` client fell more than 256 lines behind |
31+
| `console` | journald refused a VM console line |
32+
33+
Lines a `log-shim` could not journal are not counted here: the shim is a separate process and reports its drops once, when it exits.
34+
35+
Journald's internal rate-limit drops are not counted either: its notices cover a whole service and priority bucket, appear only after a later message passes the limit, and cannot provide a complete per-workload total.
36+
1937
## Gauges
2038

2139
Every gauge carries the same constant labels: `containerID`, `hostname`, `appname`, `entrypoint`, `orchestrator` and `labels` — the last being the workload's own labels, minus eru's internal ones, flattened to `k=v,k=v`.

‎logs/enc.go‎

Lines changed: 11 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -1,6 +1,7 @@
11
package logs
22

33
import (
4+
"bytes"
45
"encoding/json"
56
"io"
67
"sync"
@@ -17,19 +18,24 @@ type Encoder interface {
1718
}
1819

1920
type StreamEncoder struct {
20-
wt io.WriteCloser
21+
wt io.WriteCloser
22+
buf bytes.Buffer
23+
enc *json.Encoder
2124
}
2225

2326
func NewStreamEncoder(wt io.WriteCloser) *StreamEncoder {
24-
return &StreamEncoder{wt: wt}
27+
e := &StreamEncoder{wt: wt}
28+
e.enc = json.NewEncoder(&e.buf)
29+
return e
2530
}
2631

32+
// Encode runs under Writer's lock, which is what lets one buffer serve every line.
2733
func (e *StreamEncoder) Encode(logline *types.Log) error {
28-
data, err := json.Marshal(logline)
29-
if err != nil {
34+
e.buf.Reset()
35+
if err := e.enc.Encode(logline); err != nil {
3036
return err
3137
}
32-
_, err = e.wt.Write(append(data, '\n'))
38+
_, err := e.wt.Write(e.buf.Bytes())
3339
return err
3440
}
3541

‎logs/writer_test.go‎

Lines changed: 68 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -2,6 +2,7 @@ package logs
22

33
import (
44
"errors"
5+
"io"
56
"net"
67
"strings"
78
"testing"
@@ -127,3 +128,70 @@ func TestReconnect(t *testing.T) {
127128
writer.reconnect(ctx)
128129
assert.NoError(t, writer.Write(ctx, &types.Log{}))
129130
}
131+
132+
func BenchmarkWriterWrite(b *testing.B) {
133+
line := &types.Log{
134+
ID: "0123456789abcdef0123456789abcdef",
135+
Name: "app",
136+
Type: common.StreamStdout,
137+
EntryPoint: "web",
138+
Ident: "ident",
139+
Data: strings.Repeat("x", 200),
140+
Datetime: "2026-08-31T00:00:00.000000000Z",
141+
Extra: map[string]string{"zone": "z", "node": "n"},
142+
}
143+
ctx := b.Context()
144+
writers := map[string]*Writer{"encode": {addr: "encode", scheme: "tcp", enc: NewStreamEncoder(nopSink{})}}
145+
for name, addr := range map[string]string{"tcp": drainingTCP(b), "udp": drainingUDP(b)} {
146+
w, err := NewWriter(ctx, addr, false)
147+
require.NoError(b, err)
148+
writers[name] = w
149+
}
150+
for name, w := range writers {
151+
b.Run(name, func(b *testing.B) {
152+
b.ReportAllocs()
153+
for b.Loop() {
154+
if err := w.Write(ctx, line); err != nil {
155+
b.Fatal(err)
156+
}
157+
}
158+
})
159+
}
160+
}
161+
162+
type nopSink struct{}
163+
164+
func (nopSink) Write(p []byte) (int, error) { return len(p), nil }
165+
166+
func (nopSink) Close() error { return nil }
167+
168+
func drainingTCP(b *testing.B) string {
169+
l, err := net.Listen("tcp", "127.0.0.1:0")
170+
require.NoError(b, err)
171+
b.Cleanup(func() { _ = l.Close() })
172+
go func() {
173+
for {
174+
conn, err := l.Accept()
175+
if err != nil {
176+
return
177+
}
178+
go func() { _, _ = io.Copy(io.Discard, conn) }()
179+
}
180+
}()
181+
return "tcp://" + l.Addr().String()
182+
}
183+
184+
func drainingUDP(b *testing.B) string {
185+
conn, err := net.ListenPacket("udp", "127.0.0.1:0")
186+
require.NoError(b, err)
187+
b.Cleanup(func() { _ = conn.Close() })
188+
go func() {
189+
buf := make([]byte, 64<<10)
190+
for {
191+
if _, _, err := conn.ReadFrom(buf); err != nil {
192+
return
193+
}
194+
}
195+
}()
196+
return "udp://" + conn.LocalAddr().String()
197+
}

‎logshim/logshim.go‎

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -105,7 +105,7 @@ func (s *stream) pump(reader io.Reader) {
105105
if err != nil {
106106
return
107107
}
108-
// a journal that cannot keep up must not block the container, so a refused line is lost
108+
// a journal that refuses a line must not fail the container, so the line is lost
109109
if sendErr := s.send(string(line), s.priority, vars); sendErr != nil {
110110
s.dropped++
111111
s.err = sendErr

0 commit comments

Comments
 (0)