From d4b5dcb1fe7d849e2acbeb39d4589c4f30e48c93 Mon Sep 17 00:00:00 2001 From: gitteroy Date: Fri, 28 Aug 2026 17:26:03 +0800 Subject: [PATCH 1/6] fix(relabel): sort labels before relabel.Process for deterministic output ScratchBuilder.Labels() returns labels in Go map iteration order (random), but relabel.Process requires sorted input. Without builder.Sort(), relabel rules match non-deterministically and ~40-50% of entries reach Loki with renamed labels missing. Adds builder.Sort() plus a 1000-iteration regression test that fails without the fix. --- pkg/main.go | 4 ++- pkg/relabel_apply_test.go | 68 +++++++++++++++++++++++++++++++++++++++ 2 files changed, 71 insertions(+), 1 deletion(-) create mode 100644 pkg/relabel_apply_test.go diff --git a/pkg/main.go b/pkg/main.go index e2cdabd3..ee818a59 100644 --- a/pkg/main.go +++ b/pkg/main.go @@ -193,7 +193,9 @@ func applyRelabelConfigs(labels model.LabelSet) model.LabelSet { builder.Add(string(name), string(value)) } - // Sort labels as required by Process + // relabel.Process requires sorted input; ScratchBuilder.Labels() preserves + // insertion order (randomised map iteration), so Sort() before handing off. + builder.Sort() promLabels := builder.Labels() // Apply relabeling diff --git a/pkg/relabel_apply_test.go b/pkg/relabel_apply_test.go new file mode 100644 index 00000000..be133f1e --- /dev/null +++ b/pkg/relabel_apply_test.go @@ -0,0 +1,68 @@ +package main + +import ( + "testing" + + "github.com/prometheus/common/model" + "github.com/stretchr/testify/require" +) + +// TestApplyRelabelConfigs_DeterministicOutput verifies applyRelabelConfigs +// produces identical output across many invocations for identical input. +// +// Regression test for the upstream bug where ScratchBuilder.Labels() returns +// labels in Go map iteration order (random) and relabel.Process, which needs +// sorted input, therefore produces non-deterministic results. The fix calls +// builder.Sort() before relabel.Process. If a future upstream rebase drops +// that Sort() call, this test fails. +func TestApplyRelabelConfigs_DeterministicOutput(t *testing.T) { + defer func() { relabelConfigs = nil }() + + // Mirrors the platform default config used for AWS infra log forwarding + // (tf-common-modules/logs-forwarder relabel-configs.json). + configs, err := parseRelabelConfigs(`[ + {"source_labels":["__aws_log_type"],"target_label":"log_type","action":"replace"}, + {"source_labels":["__aws_cloudwatch_log_group"],"target_label":"log_group","action":"replace"}, + {"source_labels":["__aws_s3_log_lb"],"target_label":"loadbalancer","action":"replace"} + ]`) + require.NoError(t, err) + relabelConfigs = configs + + input := model.LabelSet{ + "__aws_log_type": "cloudwatch", + "__aws_cloudwatch_log_group": "/aws/lambda/test", + "__aws_cloudwatch_owner": "123456789012", + "aws_account_id": "123456789012", + "project_name": "test-app", + "project_env": "dev", + } + + const iters = 1000 + var firstOut model.LabelSet + for i := 0; i < iters; i++ { + in := make(model.LabelSet, len(input)) + for k, v := range input { + in[k] = v + } + out := applyRelabelConfigs(in) + + if i == 0 { + firstOut = out + require.Equal(t, model.LabelValue("cloudwatch"), out["log_type"], "log_type rename failed") + require.Equal(t, model.LabelValue("/aws/lambda/test"), out["log_group"], "log_group rename failed") + continue + } + require.Equal(t, firstOut, out, "non-deterministic output at iteration %d", i) + } +} + +// TestApplyRelabelConfigs_NoConfigsPassThrough verifies an empty config +// returns the input untouched. +func TestApplyRelabelConfigs_NoConfigsPassThrough(t *testing.T) { + defer func() { relabelConfigs = nil }() + relabelConfigs = nil + + in := model.LabelSet{"a": "1", "b": "2"} + out := applyRelabelConfigs(in) + require.Equal(t, in, out) +} From 97ee2d98261e20f10045b7d92c440ce469b4f06a Mon Sep 17 00:00:00 2001 From: gitteroy Date: Fri, 28 Aug 2026 17:28:03 +0800 Subject: [PATCH 2/6] fix(s3): skip empty batch send to prevent CloudTrail 422 DLQ failures A multi-record S3 source (e.g. CloudTrail) whose total size is a clean multiple of batchSize leaves an empty final batch. Sending an empty PushRequest makes Loki return a non-retryable 422, sending the whole SQS message to the DLQ. Guards sendToPromtail to skip zero-entry batches. Adds a behavioral test asserting no HTTP request is issued for an empty batch (fails if the guard is removed). --- pkg/empty_batch_test.go | 67 +++++++++++++++++++++++++++++++++++++++++ pkg/promtail.go | 9 +++++- 2 files changed, 75 insertions(+), 1 deletion(-) create mode 100644 pkg/empty_batch_test.go diff --git a/pkg/empty_batch_test.go b/pkg/empty_batch_test.go new file mode 100644 index 00000000..bf894246 --- /dev/null +++ b/pkg/empty_batch_test.go @@ -0,0 +1,67 @@ +package main + +import ( + "context" + "net/http" + "net/url" + "testing" + "time" + + "github.com/grafana/dskit/backoff" + "github.com/grafana/loki/v3/pkg/logproto" +) + +// recordingRoundTripper records whether any HTTP request was issued. +type recordingRoundTripper struct{ called bool } + +func (r *recordingRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) { + r.called = true + return &http.Response{StatusCode: 200, Body: http.NoBody, Header: make(http.Header)}, nil +} + +// TestSendToPromtail_EmptyBatchSkipsHTTP verifies that an empty batch never +// issues an HTTP push. Sending an empty PushRequest makes Loki return a +// non-retryable 422 that DLQs the whole SQS message. The guard in +// sendToPromtail must short-circuit before calling the HTTP client. +// +// This test FAILS if the `if entriesCount == 0 { return nil }` guard is removed. +func TestSendToPromtail_EmptyBatchSkipsHTTP(t *testing.T) { + // writeAddress is a package global read by send(); set a dummy value. + prev := writeAddress + writeAddress, _ = url.Parse("http://localhost:3100/loki/api/v1/push") + defer func() { writeAddress = prev }() + + rt := &recordingRoundTripper{} + c := &promtailClient{ + config: &promtailClientConfig{ + backoff: &backoff.Config{MinBackoff: time.Millisecond, MaxBackoff: time.Millisecond, MaxRetries: 1}, + http: &httpClientConfig{timeout: time.Second}, + }, + http: &http.Client{Transport: rt}, + log: NewLogger("error"), + } + + // Empty batch → zero entries. + b := &batch{streams: map[string]*logproto.Stream{}, processor: &LokiStages{}} + + if err := c.sendToPromtail(context.Background(), b); err != nil { + t.Fatalf("empty batch should return nil, got: %v", err) + } + if rt.called { + t.Fatal("empty batch must NOT issue an HTTP request (would trigger Loki 422 → DLQ)") + } +} + +// TestEmptyBatchReportsZeroEntries is a supporting precondition check: an empty +// batch must report zero entries so the guard above can detect it. +func TestEmptyBatchReportsZeroEntries(t *testing.T) { + b := &batch{streams: map[string]*logproto.Stream{}, processor: &LokiStages{}} + if _, cnt := b.createPushRequest(); cnt != 0 { + t.Fatalf("expected 0 entries from createPushRequest, got %d", cnt) + } + if _, entries, err := b.encode(); err != nil { + t.Fatal(err) + } else if entries != 0 { + t.Fatalf("expected 0 entries from encode, got %d", entries) + } +} diff --git a/pkg/promtail.go b/pkg/promtail.go index 122ff87a..3c39e032 100644 --- a/pkg/promtail.go +++ b/pkg/promtail.go @@ -164,11 +164,18 @@ func (b *batch) resetBatch() { } func (c *promtailClient) sendToPromtail(ctx context.Context, b *batch) error { - buf, _, err := b.encode() + buf, entriesCount, err := b.encode() if err != nil { return err } + // Skip empty batches: a multi-record source (e.g. CloudTrail) whose size is a + // clean multiple of batchSize leaves an empty final batch, and an empty + // PushRequest makes Loki return a non-retryable 422 that DLQs the whole message. + if entriesCount == 0 { + return nil + } + backoff := backoff.New(ctx, *c.config.backoff) var status int for { From 5715481d375a8b62ab201a07c3e50915f1533d42 Mon Sep 17 00:00:00 2001 From: gitteroy Date: Fri, 28 Aug 2026 17:32:12 +0800 Subject: [PATCH 3/6] ci: add fork-release workflow; disable upstream release.yml fork-release.yml builds the linux/amd64 bootstrap binary, packages a Lambda-ready zip + sha256, and publishes a GitHub release on v*-sortfix.* tags. Upstream release.yml is switched to workflow_dispatch-only because it pushes to grafanalabs S3 buckets we can't access and would otherwise fire on our v* tags. --- .github/workflows/fork-release.yml | 75 ++++++++++++++++++++++++++++++ .github/workflows/release.yml | 11 ++--- 2 files changed, 80 insertions(+), 6 deletions(-) create mode 100644 .github/workflows/fork-release.yml diff --git a/.github/workflows/fork-release.yml b/.github/workflows/fork-release.yml new file mode 100644 index 00000000..bc324005 --- /dev/null +++ b/.github/workflows/fork-release.yml @@ -0,0 +1,75 @@ +name: Fork Release + +# Builds the x86_64 Linux binary, packages it as a Lambda-ready zip, and +# publishes a GitHub release with the zip + sha256 attached when a +# v*-sortfix.* tag is pushed. +# +# Replaces upstream release.yml (which pushes to grafanalabs S3 buckets we +# don't have credentials for). Tag convention: v-sortfix., +# e.g. v1.0.1-sortfix.1. + +on: + push: + tags: + - 'v*-sortfix.*' + +permissions: + contents: write + +jobs: + release: + runs-on: ubuntu-latest + steps: + - uses: actions/checkout@v4 + + - name: Setup Go + uses: actions/setup-go@v5 + with: + go-version-file: go.mod + cache: true + + - name: Run tests + run: go test ./... + + - name: Build (linux/amd64) + run: | + GOOS=linux GOARCH=amd64 CGO_ENABLED=0 \ + go build -o bootstrap -ldflags='-s -w' ./pkg + + - name: Verify bootstrap binary + run: | + file bootstrap | grep -q 'ELF 64-bit LSB executable, x86-64' + ls -lh bootstrap + + - name: Package zip + run: | + ZIP_NAME="lambda-promtail-${GITHUB_REF_NAME}.zip" + zip "$ZIP_NAME" bootstrap + sha256sum "$ZIP_NAME" > "${ZIP_NAME}.sha256" + ls -lh "$ZIP_NAME" "${ZIP_NAME}.sha256" + echo "ZIP_NAME=$ZIP_NAME" >> "$GITHUB_ENV" + + - name: Create GitHub Release + uses: softprops/action-gh-release@v2 + with: + name: ${{ github.ref_name }} + generate_release_notes: true + files: | + ${{ env.ZIP_NAME }} + ${{ env.ZIP_NAME }}.sha256 + body: | + Fork release of `lambda-promtail` carrying two fixes for bugs still present upstream: + + - **`builder.Sort()` relabel fix** — deterministic relabel output + (see [`pkg/relabel_apply_test.go`](https://github.com/${{ github.repository }}/blob/${{ github.ref_name }}/pkg/relabel_apply_test.go)). + - **S3 empty-batch guard** — prevents CloudTrail 422 → DLQ log loss + (see [`pkg/empty_batch_test.go`](https://github.com/${{ github.repository }}/blob/${{ github.ref_name }}/pkg/empty_batch_test.go)). + + ## Artifact + + - `lambda-promtail-${{ github.ref_name }}.zip` — x86_64 Linux build for Lambda's `provided.al2023` runtime (handler: `bootstrap`). + - `lambda-promtail-${{ github.ref_name }}.zip.sha256` — checksum. + + ## Usage + + Vendored into `tf-common-modules/logs-forwarder` (`files/`) and pinned via `var.lambda_promtail_version`. diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml index 2e1a7837..ec4bf189 100644 --- a/.github/workflows/release.yml +++ b/.github/workflows/release.yml @@ -1,11 +1,10 @@ -name: Release +name: Release (upstream — disabled in fork) +# Disabled: this workflow pushes to grafanalabs S3 buckets we do not have +# credentials for. Our fork uses fork-release.yml instead. Kept for reference +# when merging upstream changes. on: - push: - branches: - - main - tags: - - 'v*' + workflow_dispatch: jobs: build: From 0956e50b11bdb59826c5d946f626aa3fa0328f9a Mon Sep 17 00:00:00 2001 From: gitteroy Date: Fri, 28 Aug 2026 17:34:54 +0800 Subject: [PATCH 4/6] test(s3): rename unused RoundTrip param to _ (revive lint) --- pkg/empty_batch_test.go | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pkg/empty_batch_test.go b/pkg/empty_batch_test.go index bf894246..7d99bba2 100644 --- a/pkg/empty_batch_test.go +++ b/pkg/empty_batch_test.go @@ -14,7 +14,7 @@ import ( // recordingRoundTripper records whether any HTTP request was issued. type recordingRoundTripper struct{ called bool } -func (r *recordingRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) { +func (r *recordingRoundTripper) RoundTrip(_ *http.Request) (*http.Response, error) { r.called = true return &http.Response{StatusCode: 200, Body: http.NoBody, Header: make(http.Header)}, nil } From 89e67cc05d013dc771da6a666bc0c4be93c48091 Mon Sep 17 00:00:00 2001 From: gitteroy Date: Sat, 29 Aug 2026 01:38:58 +0800 Subject: [PATCH 5/6] ci: add upstream-drift check that opens a rebase issue on new upstream releases --- .github/workflows/upstream-drift.yml | 82 ++++++++++++++++++++++++++++ 1 file changed, 82 insertions(+) create mode 100644 .github/workflows/upstream-drift.yml diff --git a/.github/workflows/upstream-drift.yml b/.github/workflows/upstream-drift.yml new file mode 100644 index 00000000..c1f867e0 --- /dev/null +++ b/.github/workflows/upstream-drift.yml @@ -0,0 +1,82 @@ +name: Upstream Drift Check + +# Surfaces when grafana/lambda-promtail publishes a release newer than the base +# our fork is built on. Opens a tracking issue (does NOT auto-rebase — our Sort +# and empty-batch patches can conflict, so a human does the rebase). Idempotent: +# won't reopen/duplicate an issue for a version already flagged. + +on: + schedule: + - cron: '0 1 * * 1' # Mondays 01:00 UTC + workflow_dispatch: # allow manual runs + +permissions: + contents: read + issues: write + +jobs: + check: + runs-on: ubuntu-latest + steps: + - name: Compare upstream release to our fork base + env: + GH_TOKEN: ${{ github.token }} + UPSTREAM: grafana/lambda-promtail + FORK: ${{ github.repository }} + run: | + set -euo pipefail + + # Latest upstream release tag, e.g. v1.0.2 + UPSTREAM_LATEST=$(gh release view --repo "$UPSTREAM" --json tagName --jq '.tagName') + echo "Upstream latest: $UPSTREAM_LATEST" + + # Our latest fork tag, e.g. v1.0.1-sortfix.2 → base v1.0.1 + FORK_TAG=$(gh api "repos/$FORK/tags" --jq '[.[].name | select(test("-sortfix\\."))][0]') + FORK_BASE="${FORK_TAG%%-sortfix.*}" + echo "Fork tag: $FORK_TAG (base: $FORK_BASE)" + + if [ "$UPSTREAM_LATEST" = "$FORK_BASE" ]; then + echo "Up to date — fork base matches upstream latest. Nothing to do." + exit 0 + fi + + echo "DRIFT: upstream $UPSTREAM_LATEST > fork base $FORK_BASE" + + TITLE="Upstream drift: rebase fork onto $UPSTREAM_LATEST" + + # Idempotent: skip if an open issue with this exact title already exists. + EXISTING=$(gh issue list --repo "$FORK" --state open --search "$TITLE in:title" --json number --jq 'length') + if [ "$EXISTING" != "0" ]; then + echo "Tracking issue already open — skipping." + exit 0 + fi + + # Ensure the label exists (issue create fails on a missing label). + gh label create upstream-sync --repo "$FORK" --color FBCA04 \ + --description "Upstream lambda-promtail drift" 2>/dev/null || true + + BODY=$(printf '%s\n' \ + "\`grafana/lambda-promtail\` has released **$UPSTREAM_LATEST**, newer than our fork base **$FORK_BASE** (current fork tag: \`$FORK_TAG\`)." \ + "" \ + "## Action required" \ + "Rebase our two patches onto the new upstream tag and cut a release. Both patches carry regression tests that fail if dropped — do not skip them." \ + "" \ + '```bash' \ + "git fetch upstream" \ + "git rebase $UPSTREAM_LATEST # carry Sort + empty-batch fixes forward" \ + "go test ./... # both regression tests must pass" \ + "git tag -a $UPSTREAM_LATEST-sortfix.1 -m \"$UPSTREAM_LATEST + Sort fix + S3 empty-batch fix\"" \ + "git push origin main $UPSTREAM_LATEST-sortfix.1" \ + '```' \ + "" \ + "Then vendor the new zip into the \`logs-forwarder\` module — see that module's \`UPGRADING.md\`." \ + "" \ + "If upstream has merged **both** fixes, retire the fork instead (tag the plain upstream release, drop the \`-sortfix\` suffix)." \ + "" \ + "_Opened automatically by the Upstream Drift Check workflow._") + + gh issue create --repo "$FORK" \ + --title "$TITLE" \ + --label "upstream-sync" \ + --body "$BODY" + echo "Opened tracking issue: $TITLE" From fcafef3c701449bfa6bf17123fc355b9021f951d Mon Sep 17 00:00:00 2001 From: gitteroy Date: Fri, 25 Sep 2026 16:41:54 +0800 Subject: [PATCH 6/6] feat(siem): optional S3 dual-sink for security log feed Add an S3 sink (pkg/s3_sink.go) that, when SIEM_S3_BUCKET is set, writes each batch's raw log lines (gzipped, time-partitioned) to a central SIEM bucket IN ADDITION to the Loki push. Implemented via a multiSink wrapping the existing promtailClient: Loki stays the primary/authoritative sink (governs retry/DLQ), the S3 sink is best-effort so a SIEM-side failure never blocks Loki delivery. Opt-in: no SIEM_S3_BUCKET = unchanged Loki-only behavior. Uses ambient IRSA credentials. No new dependencies (aws-sdk-go-v2 config+s3 already vendored). --- pkg/main.go | 14 ++++++ pkg/s3_sink.go | 128 +++++++++++++++++++++++++++++++++++++++++++++++++ 2 files changed, 142 insertions(+) create mode 100644 pkg/s3_sink.go diff --git a/pkg/main.go b/pkg/main.go index ee818a59..05984d98 100644 --- a/pkg/main.go +++ b/pkg/main.go @@ -250,6 +250,20 @@ func handler(ctx context.Context, ev map[string]interface{}) error { }, }, log) + // SIEM dual-sink: when SIEM_S3_BUCKET is set, ALSO write raw log lines to a + // central SIEM S3 bucket (best-effort — does not affect Loki delivery/retry). + // The Loki push (pClient) stays the primary/authoritative sink. + if siemBucket := os.Getenv("SIEM_S3_BUCKET"); siemBucket != "" { + s3sink, err := newS3Sink(ctx, siemBucket, os.Getenv("SIEM_S3_PREFIX"), log) + if err != nil { + // Do not fail ingestion if the SIEM sink can't init — log and continue + // with Loki only. + level.Error(*log).Log("msg", "siem s3 sink init failed; continuing with Loki only", "err", err) // nolint:errcheck + } else { + pClient = &multiSink{primary: pClient, secondary: []Client{s3sink}, log: log} + } + } + lokiStageConfigs, err := ParsePipelineConfigs(os.Getenv("LOKI_STAGE_CONFIGS"), *log, metrics) if err != nil { panic(err) diff --git a/pkg/s3_sink.go b/pkg/s3_sink.go new file mode 100644 index 00000000..dfc68003 --- /dev/null +++ b/pkg/s3_sink.go @@ -0,0 +1,128 @@ +package main + +import ( + "bytes" + "compress/gzip" + "context" + "crypto/rand" + "encoding/hex" + "fmt" + "strings" + "time" + + "github.com/go-kit/log" + "github.com/go-kit/log/level" + + "github.com/aws/aws-sdk-go-v2/aws" + awsconfig "github.com/aws/aws-sdk-go-v2/config" + "github.com/aws/aws-sdk-go-v2/service/s3" +) + +// s3Sink implements Client. It writes the raw log lines of a batch to a central +// SIEM S3 bucket (gzipped, time-partitioned), IN ADDITION to the Loki push done +// by promtailClient. Used for the phase-2 SIEM feed: the SIEM reads this bucket. +// +// It intentionally writes ONLY the raw log line (stream.Entries[].Line), not the +// Loki-proto wrapper, so a SIEM (OpenSearch/Data Prepper) can consume it directly. +type s3Sink struct { + client *s3.Client + bucket string + prefix string + log *log.Logger +} + +// newS3Sink builds an S3 sink using the ambient AWS credentials (IRSA role in +// the Lambda's execution environment). Region is resolved from the environment. +func newS3Sink(ctx context.Context, bucket, prefix string, logger *log.Logger) (*s3Sink, error) { + cfg, err := awsconfig.LoadDefaultConfig(ctx) + if err != nil { + return nil, fmt.Errorf("siem s3 sink: load aws config: %w", err) + } + return &s3Sink{ + client: s3.NewFromConfig(cfg), + bucket: bucket, + prefix: strings.TrimSuffix(prefix, "/"), + log: logger, + }, nil +} + +// sendToPromtail satisfies the Client interface. Despite the name (kept to match +// the interface), this writes the batch's raw log lines to S3. +func (s *s3Sink) sendToPromtail(ctx context.Context, b *batch) error { + var buf bytes.Buffer + gz := gzip.NewWriter(&buf) + + lines := 0 + for _, stream := range b.streams { + for _, e := range stream.Entries { + if e.Line == "" { + continue + } + if _, err := gz.Write([]byte(e.Line + "\n")); err != nil { + _ = gz.Close() + return fmt.Errorf("siem s3 sink: gzip write: %w", err) + } + lines++ + } + } + if err := gz.Close(); err != nil { + return fmt.Errorf("siem s3 sink: gzip close: %w", err) + } + + // Nothing to write (batch held only empty lines) — skip silently. + if lines == 0 { + return nil + } + + key := s.objectKey() + _, err := s.client.PutObject(ctx, &s3.PutObjectInput{ + Bucket: aws.String(s.bucket), + Key: aws.String(key), + Body: bytes.NewReader(buf.Bytes()), + ContentType: aws.String("application/gzip"), + ContentEncoding: aws.String("gzip"), + }) + if err != nil { + return fmt.Errorf("siem s3 sink: put object s3://%s/%s: %w", s.bucket, key, err) + } + level.Info(*s.log).Log("msg", "siem s3 sink wrote object", "bucket", s.bucket, "key", key, "lines", lines) // nolint:errcheck + return nil +} + +// objectKey builds a Hive-style time-partitioned key with a random suffix so +// concurrent Lambda invocations never collide. +func (s *s3Sink) objectKey() string { + now := time.Now().UTC() + var rnd [8]byte + _, _ = rand.Read(rnd[:]) + part := now.Format("year=2006/month=01/day=02/hour=15") + name := fmt.Sprintf("%d_%s.log.gz", now.UnixNano(), hex.EncodeToString(rnd[:])) + if s.prefix != "" { + return fmt.Sprintf("%s/%s/%s", s.prefix, part, name) + } + return fmt.Sprintf("%s/%s", part, name) +} + +// multiSink fans a batch to several Client sinks. The FIRST sink is authoritative +// for retry/error semantics (typically the Loki promtailClient): its error is +// returned so the Lambda's existing retry/DLQ behavior is unchanged. Secondary +// sinks (e.g. the SIEM S3 sink) are best-effort — their failures are logged but +// do NOT fail the invocation, so a SIEM-side problem cannot block Loki delivery. +type multiSink struct { + primary Client + secondary []Client + log *log.Logger +} + +func (m *multiSink) sendToPromtail(ctx context.Context, b *batch) error { + // Primary first — its error governs retry/DLQ. + err := m.primary.sendToPromtail(ctx, b) + + for _, sink := range m.secondary { + if serr := sink.sendToPromtail(ctx, b); serr != nil { + level.Error(*m.log).Log("msg", "secondary sink failed (best-effort, not failing invocation)", "err", serr) // nolint:errcheck + } + } + + return err +}