From e387c5616aa2731cb0b9f7ce01abe147cf8d8581 Mon Sep 17 00:00:00 2001 From: ValentaTomas Date: Sat, 16 May 2026 02:02:48 -0700 Subject: [PATCH] fix(envd): bound the in-memory logs queue Adds a byte-bounded log queue so the exporter can't grow without bound when the collector is unreachable. - 8 MiB total queue cap with drop-oldest whole entries (JSON envelopes stay intact); 192 KiB per-line ingestion cap under Loki's 256 KiB max_line_size default. - MMDS opts swapped via atomic.Pointer instead of pointer+RWMutex. - HTTP keepalives disabled on the exporter and the MMDS poll client so a pause/resume that changes endpoints doesn't reuse half-dead conns. ID injection and the send loop are unchanged; wire-compatible with the current orchestrator handler. --- packages/envd/internal/host/mmds.go | 8 +- .../envd/internal/logs/exporter/exporter.go | 77 +++++++++++++------ packages/envd/pkg/version.go | 2 +- 3 files changed, 55 insertions(+), 32 deletions(-) diff --git a/packages/envd/internal/host/mmds.go b/packages/envd/internal/host/mmds.go index 9fdc471855..a26ef588e9 100644 --- a/packages/envd/internal/host/mmds.go +++ b/packages/envd/internal/host/mmds.go @@ -38,12 +38,6 @@ type MMDSOpts struct { AccessTokenHash string `json:"accessTokenHash"` } -func (opts *MMDSOpts) Update(sandboxID, templateID, collectorAddress string) { - opts.SandboxID = sandboxID - opts.TemplateID = templateID - opts.LogsCollectorAddress = collectorAddress -} - func (opts *MMDSOpts) AddOptsToJSON(jsonLogs []byte) ([]byte, error) { parsed := make(map[string]any) @@ -136,7 +130,7 @@ func GetAccessTokenHashFromMMDS(ctx context.Context) (string, error) { } func PollForMMDSOpts(ctx context.Context, mmdsChan chan<- *MMDSOpts, envVars *utils.Map[string, string]) { - httpClient := &http.Client{} + httpClient := &http.Client{Transport: &http.Transport{DisableKeepAlives: true}} defer httpClient.CloseIdleConnections() ticker := time.NewTicker(50 * time.Millisecond) diff --git a/packages/envd/internal/logs/exporter/exporter.go b/packages/envd/internal/logs/exporter/exporter.go index da4462b148..8f273f788c 100644 --- a/packages/envd/internal/logs/exporter/exporter.go +++ b/packages/envd/internal/logs/exporter/exporter.go @@ -6,37 +6,39 @@ import ( "log" "net/http" "sync" + "sync/atomic" "time" "github.com/e2b-dev/infra/packages/envd/internal/host" ) -const ExporterTimeout = 10 * time.Second +const ( + ExporterTimeout = 10 * time.Second + + // Under Loki's 256 KiB max_line_size default. + maxLogLineBytes = 192 << 10 + maxBufferedBytes = 8 << 20 +) type HTTPExporter struct { - client http.Client - logs [][]byte - mmdsOpts *host.MMDSOpts + client http.Client + logs [][]byte + bufferedBytes int + mmdsOpts atomic.Pointer[host.MMDSOpts] // Concurrency coordination triggers chan struct{} - logLock sync.RWMutex - mmdsLock sync.RWMutex + logLock sync.Mutex startOnce sync.Once } func NewHTTPLogsExporter(ctx context.Context, mmdsChan <-chan *host.MMDSOpts) *HTTPExporter { exporter := &HTTPExporter{ client: http.Client{ - Timeout: ExporterTimeout, - }, - triggers: make(chan struct{}, 1), - startOnce: sync.Once{}, - mmdsOpts: &host.MMDSOpts{ - SandboxID: "unknown", - TemplateID: "unknown", - LogsCollectorAddress: "", + Timeout: ExporterTimeout, + Transport: &http.Transport{DisableKeepAlives: true}, }, + triggers: make(chan struct{}, 1), } go exporter.listenForMMDSOptsAndStart(ctx, mmdsChan) @@ -75,9 +77,7 @@ func (w *HTTPExporter) listenForMMDSOptsAndStart(ctx context.Context, mmdsChan < return } - w.mmdsLock.Lock() - w.mmdsOpts.Update(mmdsOpts.SandboxID, mmdsOpts.TemplateID, mmdsOpts.LogsCollectorAddress) - w.mmdsLock.Unlock() + w.mmdsOpts.Store(mmdsOpts) w.startOnce.Do(func() { go w.start(ctx) @@ -87,6 +87,11 @@ func (w *HTTPExporter) listenForMMDSOptsAndStart(ctx context.Context, mmdsChan < } func (w *HTTPExporter) start(ctx context.Context) { + // Cap stderr noise: log each error kind at most once per logFloor so a + // fast-failing collector (e.g. TCP RST) can't flood journald. + const logFloor = time.Minute + var lastLoggedJSONErr, lastLoggedSendErr time.Time + for range w.triggers { logs := w.getAllLogs() @@ -94,19 +99,27 @@ func (w *HTTPExporter) start(ctx context.Context) { continue } + opts := w.mmdsOpts.Load() + if opts == nil { + continue + } + for _, logLine := range logs { - w.mmdsLock.RLock() - logLineWithOpts, err := w.mmdsOpts.AddOptsToJSON(logLine) - w.mmdsLock.RUnlock() + logLineWithOpts, err := opts.AddOptsToJSON(logLine) if err != nil { - log.Printf("error adding instance logging options (%+v) to JSON (%+v) with logs : %v\n", w.mmdsOpts, logLine, err) + if time.Since(lastLoggedJSONErr) > logFloor { + log.Printf("error adding instance logging options to JSON: %v", err) + lastLoggedJSONErr = time.Now() + } continue } - err = w.sendInstanceLogs(ctx, logLineWithOpts, w.mmdsOpts.LogsCollectorAddress) - if err != nil { - log.Printf("error sending instance logs: %+v", err) + if err := w.sendInstanceLogs(ctx, logLineWithOpts, opts.LogsCollectorAddress); err != nil { + if time.Since(lastLoggedSendErr) > logFloor { + log.Printf("error sending instance logs: %+v", err) + lastLoggedSendErr = time.Now() + } continue } @@ -124,6 +137,11 @@ func (w *HTTPExporter) resumeProcessing() { } func (w *HTTPExporter) Write(logs []byte) (int, error) { + // Drop oversized lines: Loki would reject them anyway. + if len(logs) > maxLogLineBytes { + return len(logs), nil + } + logsCopy := make([]byte, len(logs)) copy(logsCopy, logs) @@ -138,6 +156,7 @@ func (w *HTTPExporter) getAllLogs() [][]byte { logs := w.logs w.logs = nil + w.bufferedBytes = 0 return logs } @@ -146,6 +165,16 @@ func (w *HTTPExporter) addLogs(logs []byte) { w.logLock.Lock() defer w.logLock.Unlock() + // Drop the oldest entries to stay under maxBufferedBytes. Happens when + // the collector is unreachable or the producer outruns the send loop; + // keeping the queue bounded matters more than not losing old lines. + for w.bufferedBytes+len(logs) > maxBufferedBytes && len(w.logs) > 0 { + w.bufferedBytes -= len(w.logs[0]) + w.logs[0] = nil + w.logs = w.logs[1:] + } + + w.bufferedBytes += len(logs) w.logs = append(w.logs, logs) w.resumeProcessing() diff --git a/packages/envd/pkg/version.go b/packages/envd/pkg/version.go index 27bcbdcbef..1a1c2521cc 100644 --- a/packages/envd/pkg/version.go +++ b/packages/envd/pkg/version.go @@ -1,3 +1,3 @@ package pkg -const Version = "0.5.21" +const Version = "0.5.22"