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
19 changes: 19 additions & 0 deletions .github/workflows/ci.yml
Original file line number Diff line number Diff line change
Expand Up @@ -74,6 +74,25 @@ jobs:
fail_ci_if_error: false
verbose: true

e2e:
name: E2E (durability)
runs-on: ubuntu-latest
steps:
- name: Checkout
uses: actions/checkout@df4cb1c069e1874edd31b4311f1884172cec0e10 # v6.0.3

- name: Set up Go
uses: actions/setup-go@4a3601121dd01d1626a1e23e37211e3254c1c06c # v6.4.0
with:
go-version-file: go.mod
cache: true

- name: Download dependencies
run: go mod download

- name: Run end-to-end durability tests
run: go test -tags e2e -count=1 -timeout 300s ./test/e2e/...

build:
name: Build
runs-on: ubuntu-latest
Expand Down
28 changes: 3 additions & 25 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
@@ -1,28 +1,6 @@
# Changelog

All notable changes to this project are documented here. The format is based on
[Keep a Changelog](https://keepachangelog.com/en/1.1.0/).
Release notes live on **GitHub Releases**: https://github.com/txn2/rtbeat/releases

## [Unreleased]

### Changed

- Migrated from glide/GOPATH vendoring to Go modules (`go.mod`), targeting Go 1.26.
- Upgraded Elastic libbeat to v7.17.29 (from the v7.0.0-alpha line) and refreshed all direct
dependencies (gin, prometheus client, zap).
- Replaced the libbeat-generated Makefile and Travis CI with a `make verify` workflow and GitHub
Actions (CI, CodeQL, OpenSSF Scorecard, Dependabot, docs).
- Reworked the GoReleaser pipeline to v2: static `CGO_ENABLED=0` builds for linux/darwin
(amd64, arm64), Cosign keyless signing, SBOMs, SLSA provenance, and multi-arch Docker images.

### Added

- `golangci-lint` v2 configuration and a clean lint baseline.
- `SECURITY.md`, `CODE_OF_CONDUCT.md`, `CONTRIBUTING.md`, issue/PR templates, `CODEOWNERS`,
`codecov.yml`, and MkDocs documentation.

### Notes

- The rxtx `MessageBatch` wire format is unchanged; `txn2/rxtx` is pinned to the prior production
revision and its 2018-era transitive dependencies (`coreos/bbolt`, `satori/go.uuid`) are pinned to
compatible commits.
Each tagged release (`vX.Y.Z`) ships notes generated from the commit history by
GoReleaser, alongside signed artifacts, checksums, SBOMs, and provenance.
9 changes: 9 additions & 0 deletions Makefile
Original file line number Diff line number Diff line change
Expand Up @@ -60,6 +60,14 @@ test:
@echo "==> go test (race + coverage, matches CI)"
$(GO) test -race -coverprofile=coverage.txt -covermode=atomic ./...

# End-to-end durability tests: build the real binary and drive it against a
# controllable lumberjack (logstash) server. Tag-gated so the default test run
# stays fast and offline.
.PHONY: e2e
e2e:
@echo "==> go test -tags e2e (end-to-end durability)"
$(GO) test -tags e2e -count=1 -timeout 300s ./test/e2e/...

# Install the exact golangci-lint version CI uses, into a local cache.
$(GOLANGCI_LINT):
@echo "==> installing golangci-lint $(GOLANGCI_LINT_VERSION)"
Expand Down Expand Up @@ -87,6 +95,7 @@ help:
@echo " verify lint + test + build + tidy-check + action pins"
@echo " lint golangci-lint (auto-installs $(GOLANGCI_LINT_VERSION) to .tools/)"
@echo " test go test -race with coverage profile"
@echo " e2e end-to-end durability tests (go test -tags e2e)"
@echo " build CGO_ENABLED=0 go build -o rtbeat ."
@echo " tidy go mod tidy"
@echo " tidy-check fail if go.mod/go.sum are not tidy"
Expand Down
158 changes: 128 additions & 30 deletions beater/rtbeat.go
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,8 @@ import (
"io"
"log"
"net/http"
"sync"
"sync/atomic"
"time"

"github.com/prometheus/client_golang/prometheus/promhttp"
Expand Down Expand Up @@ -45,6 +47,16 @@ func New(_ *beat.Beat, cfg *common.Config) (beat.Beater, error) {
return nil, fmt.Errorf("error reading config file: %v", err)
}

// Guard the durability footguns: a non-positive ack timeout would 504 every
// request, and a non-positive shutdown timeout would silently disable the
// drain. Fall back to the defaults.
if c.Timeout <= 0 {
c.Timeout = config.DefaultConfig.Timeout
}
if c.ShutdownTimeout <= 0 {
c.ShutdownTimeout = config.DefaultConfig.ShutdownTimeout
}

zapCfg := zap.NewProductionConfig()
zapCfg.DisableCaller = true
zapCfg.DisableStacktrace = true
Expand Down Expand Up @@ -98,16 +110,30 @@ func (bt *Rtbeat) Run(b *beat.Beat) error {

var err error
bt.client, err = b.Publisher.ConnectWith(beat.ClientConfig{
//PublishMode: beat.GuaranteedSend,
ACKHandler: acker.RawCounting(func(i int) {
bt.logger.Info("Run", zapcore.Field{
Key: "ACKCount",
Type: zapcore.Int32Type,
Integer: int64(i),
})
currentAcks.Set(float64(i))
totalAcks.Add(float64(i))
}),
// GuaranteedSend retries events until the output acknowledges them;
// WaitClose makes Close() (called from Stop) block until in-flight
// events are acked or the shutdown budget elapses, so a graceful
// shutdown drains rather than dropping in-flight events.
PublishMode: beat.GuaranteedSend,
WaitClose: time.Duration(bt.config.ShutdownTimeout) * time.Second,
ACKHandler: acker.Combine(
// Existing ack metrics.
acker.RawCounting(func(i int) {
bt.logger.Info("Run", zapcore.Field{
Key: "ACKCount",
Type: zapcore.Int32Type,
Integer: int64(i),
})
currentAcks.Set(float64(i))
totalAcks.Add(float64(i))
}),
// Per-batch ack correlation: each event carries its *batchAck in
// Private; when all of a batch's events are acked, its waiter is
// released so POST /in can return 200 (delivered).
acker.EventPrivateReporter(func(_ int, data []interface{}) {
resolveBatchAcks(data)
}),
),
})
if err != nil {
return err
Expand All @@ -123,7 +149,8 @@ func (bt *Rtbeat) Run(b *beat.Beat) error {
// get a router
r := gin.Default()

r.POST("/in", inHandler(b.Info.Name, bt.logger, bt.client, batches.Inc, func(n int) {
ackTimeout := time.Duration(bt.config.Timeout) * time.Second
r.POST("/in", inHandler(b.Info.Name, bt.logger, bt.client, ackTimeout, batches.Inc, func(n int) {
messages.Add(float64(n))
}))

Expand Down Expand Up @@ -164,16 +191,25 @@ func (bt *Rtbeat) Run(b *beat.Beat) error {
},
)

// shutdown the web server
ctx, cancel := context.WithTimeout(context.Background(), 5*time.Second)
// Graceful shutdown ordering matters for durability: stop accepting new
// requests first (Shutdown waits for in-flight handlers to finish waiting
// on their acks), THEN close the client. Closing was configured with
// WaitClose + GuaranteedSend, so Close drains outstanding events until
// acked or the shutdown budget elapses. Closing before intake stops would
// let a late batch publish into a closing pipeline and be lost.
shutdownTimeout := time.Duration(bt.config.ShutdownTimeout) * time.Second
ctx, cancel := context.WithTimeout(context.Background(), shutdownTimeout)
_ = srv.Shutdown(ctx)
cancel()
_ = bt.client.Close()
return nil
}

// Stop the beat
// Stop signals the beat to shut down. The actual drain (stop intake, finish
// in-flight handlers, then drain and close the publisher) happens in Run's
// shutdown sequence to guarantee the ordering.
func (bt *Rtbeat) Stop() {
_ = bt.client.Close()
bt.logger.Info("Stop", zap.String("state", "shutting down"))
close(bt.done)
}

Expand All @@ -183,32 +219,73 @@ type eventPublisher interface {
PublishAll([]beat.Event)
}

// batchAck tracks the outstanding acknowledgements for a single published
// batch. done is closed once every event in the batch has been acked by the
// output, releasing the HTTP handler waiting on delivery.
type batchAck struct {
remaining int64
done chan struct{}
once sync.Once
}

func newBatchAck(n int) *batchAck {
return &batchAck{remaining: int64(n), done: make(chan struct{})}
}

// ack records n acknowledged events and releases waiters once the batch is
// fully delivered.
func (a *batchAck) ack(n int) {
if atomic.AddInt64(&a.remaining, -int64(n)) <= 0 {
a.once.Do(func() { close(a.done) })
}
}

// resolveBatchAcks groups a slice of acked events' Private values by their
// owning batch and releases each batch that is now fully delivered. The output
// may ack events from several batches in a single callback, and may split a
// batch's acks across callbacks; grouping per *batchAck handles both. Non
// *batchAck / nil entries are ignored.
func resolveBatchAcks(data []interface{}) {
counts := make(map[*batchAck]int, len(data))
for _, d := range data {
if a, ok := d.(*batchAck); ok && a != nil {
counts[a]++
}
}
for a, n := range counts {
a.ack(n)
}
}

// buildEvents converts an rxtx MessageBatch into the beat.Event slice rtbeat
// publishes. Each message becomes one event carrying the original message under
// "rxtxMsg" alongside "type" and "clientIp". The slice is pre-sized with one
// leading zero-value event; this is long-standing behavior, preserved here
// intentionally (changing it is tracked separately).
func buildEvents(beatName, clientIP string, msg *rtq.MessageBatch) []beat.Event {
events := make([]beat.Event, 1)
for i, message := range msg.Messages {
// publishes — one event per message carrying the original message under
// "rxtxMsg" alongside "type" and "clientIp". Each event's Private holds the
// shared *batchAck so the ACK handler can correlate acks back to this batch.
func buildEvents(beatName, clientIP string, msg *rtq.MessageBatch, ack *batchAck) []beat.Event {
events := make([]beat.Event, 0, len(msg.Messages))
for _, message := range msg.Messages {
events = append(events, beat.Event{
Timestamp: time.Now(),
Fields: common.MapStr{
"type": beatName,
"rxtxMsg": message,
"clientIp": clientIP,
},
Private: i,
Private: ack,
})
}
return events
}

// inHandler builds the POST /in gin handler. Metric updates are injected as
// callbacks (onBatch, onMessages) so the handler can be exercised in tests
// without registering against the global prometheus registry. The handler
// responds before publishing so a slow output never blocks the rxtx client.
func inHandler(beatName string, logger *zap.Logger, pub eventPublisher, onBatch func(), onMessages func(n int)) gin.HandlerFunc {
// without registering against the global prometheus registry.
//
// Durability: the handler publishes the batch and then waits for the output to
// acknowledge delivery before responding 200, bounded by ackTimeout. If the
// ack does not arrive in time it responds 504 so the sender (e.g. rxtx) keeps
// its durable copy and retries, rather than dropping it on a premature 200.
func inHandler(beatName string, logger *zap.Logger, pub eventPublisher, ackTimeout time.Duration, onBatch func(), onMessages func(n int)) gin.HandlerFunc {
return func(c *gin.Context) {
onBatch()

Expand All @@ -224,14 +301,35 @@ func inHandler(beatName string, logger *zap.Logger, pub eventPublisher, onBatch
return
}

// respond quickly to avoid getting a re-send from the server
c.JSON(http.StatusOK, gin.H{"status": "OK"})
n := len(msg.Messages)
onMessages(n)

// Nothing to deliver: acknowledge immediately.
if n == 0 {
c.JSON(http.StatusOK, gin.H{"status": "OK"})
return
}

onMessages(len(msg.Messages))
// Capture the client IP synchronously: the gin context is recycled
// once the handler returns, so it must not be read from the goroutine.
events := buildEvents(beatName, c.ClientIP(), msg)
ack := newBatchAck(n)
events := buildEvents(beatName, c.ClientIP(), msg, ack)

go pub.PublishAll(events)

select {
case <-ack.done:
c.JSON(http.StatusOK, gin.H{"status": "OK"})
case <-time.After(ackTimeout):
logger.Warn("Run",
zap.String("state", "ack timeout"),
zap.Int("messages", n),
zap.Duration("timeout", ackTimeout),
)
c.JSON(http.StatusGatewayTimeout, gin.H{
"status": "TIMEOUT",
"message": "events accepted but not acknowledged by the output within timeout",
})
}
}
}
Loading