Skip to content
Closed
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
75 changes: 75 additions & 0 deletions .github/workflows/fork-release.yml
Original file line number Diff line number Diff line change
@@ -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<upstream-base>-sortfix.<n>,
# 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`.
11 changes: 5 additions & 6 deletions .github/workflows/release.yml
Original file line number Diff line number Diff line change
@@ -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:
Expand Down
82 changes: 82 additions & 0 deletions .github/workflows/upstream-drift.yml
Original file line number Diff line number Diff line change
@@ -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"
67 changes: 67 additions & 0 deletions pkg/empty_batch_test.go
Original file line number Diff line number Diff line change
@@ -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(_ *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)
}
}
18 changes: 17 additions & 1 deletion pkg/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -248,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)
Expand Down
9 changes: 8 additions & 1 deletion pkg/promtail.go
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
68 changes: 68 additions & 0 deletions pkg/relabel_apply_test.go
Original file line number Diff line number Diff line change
@@ -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)
}
Loading
Loading