diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index ea2fc89b..fbcb811d 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -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 diff --git a/CHANGELOG.md b/CHANGELOG.md index 26904f58..3b08801b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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. diff --git a/Makefile b/Makefile index 84ebfc9d..0b6107be 100644 --- a/Makefile +++ b/Makefile @@ -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)" @@ -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" diff --git a/beater/rtbeat.go b/beater/rtbeat.go index 1dc492a9..325398c2 100644 --- a/beater/rtbeat.go +++ b/beater/rtbeat.go @@ -7,6 +7,8 @@ import ( "io" "log" "net/http" + "sync" + "sync/atomic" "time" "github.com/prometheus/client_golang/prometheus/promhttp" @@ -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 @@ -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 @@ -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)) })) @@ -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) } @@ -183,14 +219,51 @@ 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{ @@ -198,7 +271,7 @@ func buildEvents(beatName, clientIP string, msg *rtq.MessageBatch) []beat.Event "rxtxMsg": message, "clientIp": clientIP, }, - Private: i, + Private: ack, }) } return events @@ -206,9 +279,13 @@ func buildEvents(beatName, clientIP string, msg *rtq.MessageBatch) []beat.Event // 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() @@ -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", + }) + } } } diff --git a/beater/rtbeat_test.go b/beater/rtbeat_test.go index f45545b1..b92d9af4 100644 --- a/beater/rtbeat_test.go +++ b/beater/rtbeat_test.go @@ -16,16 +16,26 @@ import ( func init() { gin.SetMode(gin.TestMode) } // fakePublisher records the events handed to PublishAll, satisfying the -// eventPublisher interface in place of a real libbeat client. +// eventPublisher interface in place of a real libbeat client. When autoAck is +// set it simulates the output delivering the batch by acking the batchAck +// carried in the events' Private field. type fakePublisher struct { - got chan []beat.Event + got chan []beat.Event + autoAck bool } -func newFakePublisher() *fakePublisher { - return &fakePublisher{got: make(chan []beat.Event, 1)} +func newFakePublisher(autoAck bool) *fakePublisher { + return &fakePublisher{got: make(chan []beat.Event, 1), autoAck: autoAck} } -func (f *fakePublisher) PublishAll(events []beat.Event) { f.got <- events } +func (f *fakePublisher) PublishAll(events []beat.Event) { + f.got <- events + if f.autoAck && len(events) > 0 { + if a, ok := events[0].Private.(*batchAck); ok && a != nil { + a.ack(len(events)) + } + } +} func TestBuildEvents(t *testing.T) { msg := &rtq.MessageBatch{ @@ -36,20 +46,17 @@ func TestBuildEvents(t *testing.T) { {Seq: "2", Producer: "p", Payload: map[string]interface{}{"b": 2.0}}, }, } + ack := newBatchAck(len(msg.Messages)) - events := buildEvents("rtbeat", "10.0.0.1", msg) + events := buildEvents("rtbeat", "10.0.0.1", msg, ack) - // One leading zero-value placeholder event (preserved legacy behavior), - // then one event per message. - if got, want := len(events), 1+len(msg.Messages); got != want { + // One event per message — no placeholder. + if got, want := len(events), len(msg.Messages); got != want { t.Fatalf("len(events) = %d, want %d", got, want) } - if events[0].Fields != nil { - t.Errorf("events[0] should be the zero-value placeholder, got Fields=%v", events[0].Fields) - } for i, in := range msg.Messages { - ev := events[i+1] + ev := events[i] if ev.Fields["type"] != "rtbeat" { t.Errorf("event %d type = %v, want rtbeat", i, ev.Fields["type"]) } @@ -63,8 +70,9 @@ func TestBuildEvents(t *testing.T) { if got.Seq != in.Seq { t.Errorf("event %d rxtxMsg.Seq = %q, want %q", i, got.Seq, in.Seq) } - if ev.Private != i { - t.Errorf("event %d Private = %v, want %d", i, ev.Private, i) + // Every event carries the shared batch ack token in Private. + if ev.Private != ack { + t.Errorf("event %d Private = %v, want the batchAck", i, ev.Private) } if ev.Timestamp.IsZero() { t.Errorf("event %d Timestamp is zero", i) @@ -72,23 +80,80 @@ func TestBuildEvents(t *testing.T) { } } +func isClosed(ch <-chan struct{}) bool { + select { + case <-ch: + return true + default: + return false + } +} + +func TestBatchAckPartial(t *testing.T) { + a := newBatchAck(3) + + a.ack(2) + if isClosed(a.done) { + t.Fatal("done closed after 2 of 3 acks") + } + + a.ack(1) + if !isClosed(a.done) { + t.Fatal("done not closed after all 3 acks") + } + + // Extra acks must not panic or re-close. + a.ack(1) +} + +func TestResolveBatchAcks(t *testing.T) { + // Two batches whose events are interleaved in a single ack callback, + // plus an unrelated/nil entry that must be ignored. + a := newBatchAck(2) + b := newBatchAck(1) + + resolveBatchAcks([]interface{}{a, b, a, nil, 42}) + + if !isClosed(a.done) { + t.Error("batch a (2 events, both acked) should be done") + } + if !isClosed(b.done) { + t.Error("batch b (1 event, acked) should be done") + } +} + +func TestResolveBatchAcksSplitAcrossCallbacks(t *testing.T) { + // A 3-event batch acked as 2 then 1 across two reporter callbacks. + a := newBatchAck(3) + + resolveBatchAcks([]interface{}{a, a}) + if isClosed(a.done) { + t.Fatal("batch closed after only 2 of 3 acks") + } + + resolveBatchAcks([]interface{}{a}) + if !isClosed(a.done) { + t.Fatal("batch not closed after the remaining ack") + } +} + func TestBuildEventsEmptyBatch(t *testing.T) { - events := buildEvents("rtbeat", "1.2.3.4", &rtq.MessageBatch{}) - if len(events) != 1 { - t.Fatalf("empty batch: len(events) = %d, want 1 (placeholder only)", len(events)) + events := buildEvents("rtbeat", "1.2.3.4", &rtq.MessageBatch{}, newBatchAck(0)) + if len(events) != 0 { + t.Fatalf("empty batch: len(events) = %d, want 0", len(events)) } } -func newTestRouter(pub eventPublisher, onBatch func(), onMessages func(int)) *gin.Engine { +func newTestRouter(pub eventPublisher, ackTimeout time.Duration, onBatch func(), onMessages func(int)) *gin.Engine { r := gin.New() - r.POST("/in", inHandler("rtbeat", zap.NewNop(), pub, onBatch, onMessages)) + r.POST("/in", inHandler("rtbeat", zap.NewNop(), pub, ackTimeout, onBatch, onMessages)) return r } -func TestInHandlerValidBatch(t *testing.T) { - pub := newFakePublisher() +func TestInHandlerDeliveredBatch(t *testing.T) { + pub := newFakePublisher(true) // simulate the output acking the batch var batches, messages int - r := newTestRouter(pub, func() { batches++ }, func(n int) { messages += n }) + r := newTestRouter(pub, 2*time.Second, func() { batches++ }, func(n int) { messages += n }) body := `{"uuid":"b1","size":1,"messages":[{"seq":"1","payload":{"hello":"world"}}]}` w := httptest.NewRecorder() @@ -108,12 +173,10 @@ func TestInHandlerValidBatch(t *testing.T) { select { case events := <-pub.got: - if len(events) != 2 { // placeholder + one message - t.Fatalf("published %d events, want 2", len(events)) + if len(events) != 1 { // one message, no placeholder + t.Fatalf("published %d events, want 1", len(events)) } - // clientIp is read synchronously from the request; assert it - // propagates through the handler into the published event. - if got := events[1].Fields["clientIp"]; got != "203.0.113.7" { + if got := events[0].Fields["clientIp"]; got != "203.0.113.7" { t.Errorf("published clientIp = %v, want 203.0.113.7", got) } case <-time.After(2 * time.Second): @@ -121,10 +184,31 @@ func TestInHandlerValidBatch(t *testing.T) { } } +func TestInHandlerAckTimeout(t *testing.T) { + pub := newFakePublisher(false) // output never acks + r := newTestRouter(pub, 50*time.Millisecond, func() {}, func(int) {}) + + body := `{"uuid":"b1","size":1,"messages":[{"seq":"1","payload":{"hello":"world"}}]}` + w := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/in", bytes.NewBufferString(body)) + r.ServeHTTP(w, req) + + if w.Code != http.StatusGatewayTimeout { + t.Fatalf("status = %d, want 504; body=%s", w.Code, w.Body.String()) + } + // The batch was still published (it may yet be delivered; the sender + // retries on the 504). + select { + case <-pub.got: + case <-time.After(time.Second): + t.Fatal("PublishAll was not called") + } +} + func TestInHandlerBadJSON(t *testing.T) { - pub := newFakePublisher() + pub := newFakePublisher(true) var batches, messages int - r := newTestRouter(pub, func() { batches++ }, func(n int) { messages += n }) + r := newTestRouter(pub, 2*time.Second, func() { batches++ }, func(n int) { messages += n }) w := httptest.NewRecorder() req := httptest.NewRequest(http.MethodPost, "/in", bytes.NewBufferString(`{not valid json`)) @@ -140,9 +224,29 @@ func TestInHandlerBadJSON(t *testing.T) { t.Errorf("onMessages total = %d, want 0 on parse failure", messages) } - // On the error path the handler returns before spawning the publish - // goroutine, so nothing can ever be sent — assert deterministically. + // On the error path the handler returns before publishing. if len(pub.got) != 0 { t.Fatal("PublishAll must not be called on bad JSON") } } + +func TestInHandlerEmptyBatch(t *testing.T) { + pub := newFakePublisher(true) + var messages int + r := newTestRouter(pub, 2*time.Second, func() {}, func(n int) { messages += n }) + + w := httptest.NewRecorder() + req := httptest.NewRequest(http.MethodPost, "/in", bytes.NewBufferString(`{"uuid":"b","size":0,"messages":[]}`)) + r.ServeHTTP(w, req) + + if w.Code != http.StatusOK { + t.Fatalf("status = %d, want 200 for empty batch", w.Code) + } + if messages != 0 { + t.Errorf("onMessages total = %d, want 0", messages) + } + // Nothing to publish for an empty batch. + if len(pub.got) != 0 { + t.Fatal("PublishAll must not be called for an empty batch") + } +} diff --git a/beater/run_test.go b/beater/run_test.go index 9b4abb74..afd81a8e 100644 --- a/beater/run_test.go +++ b/beater/run_test.go @@ -14,22 +14,36 @@ import ( ) // fakeClient implements beat.Client, recording published events and signaling -// when it is closed. +// when it is closed. It simulates the output delivering events by driving the +// registered ACK handler, which is how rtbeat's per-batch ack waiter and ack +// metrics are exercised end-to-end. type fakeClient struct { published chan []beat.Event closed chan struct{} + acker beat.ACKer } func newFakeClient() *fakeClient { return &fakeClient{published: make(chan []beat.Event, 8), closed: make(chan struct{})} } -func (c *fakeClient) Publish(e beat.Event) { c.published <- []beat.Event{e} } -func (c *fakeClient) PublishAll(es []beat.Event) { c.published <- es } -func (c *fakeClient) Close() error { close(c.closed); return nil } +func (c *fakeClient) Publish(e beat.Event) { c.PublishAll([]beat.Event{e}) } + +func (c *fakeClient) PublishAll(es []beat.Event) { + c.published <- es + if c.acker != nil { + for _, e := range es { + c.acker.AddEvent(e, true) + } + c.acker.ACKEvents(len(es)) + } +} + +func (c *fakeClient) Close() error { close(c.closed); return nil } // fakePipeline implements beat.Pipeline and captures the ClientConfig so the -// test can drive the ACK handler. +// test can inspect it; it also hands the client the registered ACK handler so +// publishes can simulate delivery. type fakePipeline struct { client *fakeClient cfg beat.ClientConfig @@ -37,6 +51,7 @@ type fakePipeline struct { func (p *fakePipeline) ConnectWith(cc beat.ClientConfig) (beat.Client, error) { p.cfg = cc + p.client.acker = cc.ACKHandler return p.client, nil } @@ -72,6 +87,28 @@ func waitForServer(t *testing.T, url string) { t.Fatal("server did not start within timeout") } +func TestNewClampsNonPositiveTimeouts(t *testing.T) { + cfg, err := common.NewConfigFrom(map[string]interface{}{ + "timeout": 0, + "shutdown_timeout": -1, + }) + if err != nil { + t.Fatalf("config: %v", err) + } + bter, err := New(nil, cfg) + if err != nil { + t.Fatalf("New: %v", err) + } + bt := bter.(*Rtbeat) + + if bt.config.Timeout != 5 { + t.Errorf("Timeout = %d, want clamped to default 5", bt.config.Timeout) + } + if bt.config.ShutdownTimeout != 30 { + t.Errorf("ShutdownTimeout = %d, want clamped to default 30", bt.config.ShutdownTimeout) + } +} + // TestRunLifecycle drives the full beat through a fake libbeat pipeline: // New -> Run -> serve /in and /metrics -> ACK -> Stop. func TestRunLifecycle(t *testing.T) { @@ -109,7 +146,9 @@ func TestRunLifecycle(t *testing.T) { base := "http://127.0.0.1:" + port waitForServer(t, base+"/metrics") - // POST a batch and confirm it is published through the pipeline. + // POST a batch. The fake client acks it through the registered ACK + // handler, so the durability path (publish -> wait for delivery -> 200) + // completes and returns 200. body := `{"uuid":"b1","size":1,"messages":[{"seq":"1","payload":{"k":"v"}}]}` resp, err := http.Post(base+"/in", "application/json", strings.NewReader(body)) //nolint:gosec // fixed localhost test URL if err != nil { @@ -123,8 +162,8 @@ func TestRunLifecycle(t *testing.T) { select { case events := <-client.published: - if len(events) != 2 { // placeholder + one message - t.Errorf("published %d events, want 2", len(events)) + if len(events) != 1 { // one message, no placeholder + t.Errorf("published %d events, want 1", len(events)) } case <-time.After(2 * time.Second): t.Fatal("pipeline did not receive published events") @@ -141,13 +180,12 @@ func TestRunLifecycle(t *testing.T) { t.Errorf("GET /metrics status = %d, want 200", mResp.StatusCode) } - // Exercise the ACK callback wired up in Run. + // The ACK handler must have been registered for delivery confirmation. if pipe.cfg.ACKHandler == nil { t.Fatal("Run did not register an ACK handler") } - pipe.cfg.ACKHandler.ACKEvents(3) - // Stop closes the client and unblocks Run. + // Stop closes the client (draining) and unblocks Run. bt.Stop() select { diff --git a/config/config.go b/config/config.go index 0d633178..14ec0d40 100644 --- a/config/config.go +++ b/config/config.go @@ -4,11 +4,18 @@ package config type Config struct { - Port string `config:"port"` - Timeout int `config:"timeout"` + Port string `config:"port"` + // Timeout is the per-request budget, in seconds, to wait for the output + // to acknowledge a batch before POST /in responds 504. Durability hinges + // on this: the 200 means "delivered", not merely "received". + Timeout int `config:"timeout"` + // ShutdownTimeout is the maximum time, in seconds, Stop() waits for + // in-flight events to be acknowledged before closing the publisher. + ShutdownTimeout int `config:"shutdown_timeout"` } var DefaultConfig = Config{ - Port: "8081", - Timeout: 5, + Port: "8081", + Timeout: 5, + ShutdownTimeout: 30, } diff --git a/config/config_test.go b/config/config_test.go index 07a7e16d..c7f1082a 100644 --- a/config/config_test.go +++ b/config/config_test.go @@ -13,12 +13,16 @@ func TestDefaultConfig(t *testing.T) { if DefaultConfig.Timeout != 5 { t.Errorf("default Timeout = %d, want 5", DefaultConfig.Timeout) } + if DefaultConfig.ShutdownTimeout != 30 { + t.Errorf("default ShutdownTimeout = %d, want 30", DefaultConfig.ShutdownTimeout) + } } func TestUnpackOverrides(t *testing.T) { raw, err := common.NewConfigFrom(map[string]interface{}{ - "port": "9090", - "timeout": 30, + "port": "9090", + "timeout": 30, + "shutdown_timeout": 45, }) if err != nil { t.Fatalf("NewConfigFrom: %v", err) @@ -35,6 +39,9 @@ func TestUnpackOverrides(t *testing.T) { if c.Timeout != 30 { t.Errorf("Timeout = %d, want 30", c.Timeout) } + if c.ShutdownTimeout != 45 { + t.Errorf("ShutdownTimeout = %d, want 45", c.ShutdownTimeout) + } } func TestUnpackEmptyKeepsDefaults(t *testing.T) { diff --git a/docs/configuration.md b/docs/configuration.md index f6b90357..8556c446 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -9,11 +9,28 @@ processor, and logging options is documented in `rtbeat.reference.yml` at the re rtbeat: # HTTP port the /in and /metrics endpoints listen on. port: "8081" - # Read timeout, in seconds. + # Seconds POST /in waits for the output to acknowledge a batch (see Delivery). timeout: 5 + # Seconds Stop() waits to drain in-flight events on graceful shutdown. + shutdown_timeout: 30 ``` -The defaults (port `8081`, timeout `5`) are applied when the keys are omitted. +The defaults (port `8081`, timeout `5`, shutdown_timeout `30`) are applied when the keys are omitted. + +## Delivery semantics + +rtbeat acknowledges a batch only **after** the configured output confirms delivery: + +- `POST /in` publishes the batch with `GuaranteedSend` and waits up to `timeout` seconds for the + output to ack. On ack it returns **200** ("delivered"); if the ack does not arrive in time it + returns **504**, so a durable sender such as [rxtx](https://github.com/txn2/rxtx) keeps its copy + and retries rather than dropping it on a premature 200. (Use a deterministic document id downstream + so a retried batch dedupes.) +- On graceful shutdown, rtbeat stops accepting new requests, lets in-flight handlers finish, then + drains outstanding events for up to `shutdown_timeout` seconds before closing the publisher. + +For at-least-once delivery across hard crashes as well, configure libbeat's disk queue/spool in the +output/queue settings (see `rtbeat.reference.yml`). ## Output diff --git a/go.mod b/go.mod index 63fb9ebe..366f5ccf 100644 --- a/go.mod +++ b/go.mod @@ -29,6 +29,7 @@ replace ( require ( github.com/elastic/beats/v7 v7.17.29 + github.com/elastic/go-lumber v0.1.0 github.com/gin-gonic/gin v1.12.0 github.com/prometheus/client_golang v1.23.2 github.com/txn2/rxtx v1.5.2 @@ -64,7 +65,6 @@ require ( github.com/elastic/elastic-agent-client/v7 v7.8.1 // indirect github.com/elastic/elastic-agent-libs v0.7.2 // indirect github.com/elastic/go-concert v0.2.0 // indirect - github.com/elastic/go-lumber v0.1.0 // indirect github.com/elastic/go-seccomp-bpf v1.2.0 // indirect github.com/elastic/go-structform v0.0.9 // indirect github.com/elastic/go-sysinfo v1.8.1 // indirect diff --git a/rtbeat.yml b/rtbeat.yml index 446e0967..5670c2e7 100644 --- a/rtbeat.yml +++ b/rtbeat.yml @@ -5,8 +5,13 @@ rtbeat: # HTTP server port port: "8081" - # timeout + # Seconds POST /in waits for the output to acknowledge a batch before + # responding 504. The 200 means "delivered", not merely "received", so the + # sender (e.g. rxtx) keeps its durable copy until delivery is confirmed. timeout: 5 + # Seconds Stop() waits for in-flight events to be acknowledged (drained) + # before closing the publisher on graceful shutdown. + shutdown_timeout: 30 #================================ General ===================================== diff --git a/test/e2e/durability_test.go b/test/e2e/durability_test.go new file mode 100644 index 00000000..c6389aa3 --- /dev/null +++ b/test/e2e/durability_test.go @@ -0,0 +1,422 @@ +//go:build e2e + +// Package e2e black-box tests the rtbeat binary's durability guarantees against +// a real libbeat logstash output talking to a controllable lumberjack server. +// +// Run with: go test -tags e2e ./test/e2e/... +// +// It proves, end to end (real binary, real TCP, real lumberjack ack protocol, +// real SIGTERM path), that: +// - POST /in returns 200 only AFTER the output acknowledges the batch (and +// not before — see TestE2EWaitsForAckBeforeOK, which delays the ack); +// - POST /in returns 504 when the output never acks (so the sender retries); +// - a single /in batch is acked correctly even when the output fragments the +// acks across the wire (TestE2EMultiMessageFragmentedAck); +// - an in-flight request is still delivered when SIGTERM arrives mid-flight +// (graceful drain), and the process exits cleanly. +// +// Out of scope (NOT proven here, by design): durability across a hard crash +// (SIGKILL/panic/power loss) and across process restarts. rtbeat holds +// in-flight events in libbeat's in-memory queue; crash-safety relies on the +// sender retrying on a non-2xx response and, optionally, libbeat's disk +// queue/spool — neither is exercised by this suite. +package e2e + +import ( + "errors" + "fmt" + "io" + "net" + "net/http" + "os" + "os/exec" + "path/filepath" + "strings" + "sync" + "syscall" + "testing" + "time" + + "github.com/elastic/go-lumber/lj" + "github.com/elastic/go-lumber/server" +) + +const validBatch = `{"uuid":"e2e","size":1,"messages":[{"seq":"1","producer":"test","payload":{"k":"v"}}]}` + +var rtbeatBin string + +func TestMain(m *testing.M) { + root, err := repoRoot() + if err != nil { + fmt.Fprintln(os.Stderr, "e2e: locate repo root:", err) + os.Exit(1) + } + + bin := filepath.Join(os.TempDir(), fmt.Sprintf("rtbeat-e2e-%d", os.Getpid())) + build := exec.Command("go", "build", "-o", bin, ".") + build.Dir = root + build.Env = append(os.Environ(), "CGO_ENABLED=0") + if out, err := build.CombinedOutput(); err != nil { + fmt.Fprintf(os.Stderr, "e2e: build rtbeat: %v\n%s\n", err, out) + os.Exit(1) + } + rtbeatBin = bin + + code := m.Run() + _ = os.Remove(bin) + os.Exit(code) +} + +func repoRoot() (string, error) { + dir, err := os.Getwd() + if err != nil { + return "", err + } + for { + if _, err := os.Stat(filepath.Join(dir, "go.mod")); err == nil { + return dir, nil + } + parent := filepath.Dir(dir) + if parent == dir { + return "", errors.New("go.mod not found above test directory") + } + dir = parent + } +} + +// ackPolicy controls how the lumberjack server acknowledges received batches. +type ackPolicy int + +const ( + ackImmediate ackPolicy = iota // ack as soon as received -> 200 + ackNever // never ack -> handler times out -> 504 + ackDelayed // ack after a delay -> exercises ordering / drain +) + +// lumberServer is a controllable lumberjack (logstash) endpoint. +type lumberServer struct { + srv server.Server + addr string + mu sync.Mutex + events int +} + +func startLumber(t *testing.T, policy ackPolicy, delay time.Duration) *lumberServer { + t.Helper() + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("lumber listen: %v", err) + } + s, err := server.NewWithListener(l, server.V2(true)) + if err != nil { + _ = l.Close() + t.Fatalf("lumber server: %v", err) + } + ls := &lumberServer{srv: s, addr: l.Addr().String()} + + go func() { + for batch := range s.ReceiveChan() { + ls.mu.Lock() + ls.events += len(batch.Events) + ls.mu.Unlock() + + switch policy { + case ackImmediate: + batch.ACK() + case ackDelayed: + go func(b *lj.Batch) { + time.Sleep(delay) + b.ACK() + }(batch) + case ackNever: + // hold the ack to force the rtbeat handler to time out + } + } + }() + + return ls +} + +func (ls *lumberServer) received() int { + ls.mu.Lock() + defer ls.mu.Unlock() + return ls.events +} + +func (ls *lumberServer) close() { _ = ls.srv.Close() } + +// rtbeatProc is a running rtbeat binary under test. A single goroutine owns +// cmd.Wait (its result is delivered on waitErr) so signalling and killing never +// race on a second Wait. +type rtbeatProc struct { + cmd *exec.Cmd + httpPort int + stderr *syncBuffer + waitErr chan error +} + +func startRtbeat(t *testing.T, lsAddr string, timeoutSec int) *rtbeatProc { + return startRtbeatCfg(t, lsAddr, timeoutSec, 0) +} + +func startRtbeatCfg(t *testing.T, lsAddr string, timeoutSec, bulkMaxSize int) *rtbeatProc { + t.Helper() + httpPort := freePort(t) + + bulk := "" + if bulkMaxSize > 0 { + bulk = fmt.Sprintf("\n bulk_max_size: %d", bulkMaxSize) + } + cfg := fmt.Sprintf(`rtbeat: + port: "%d" + timeout: %d + shutdown_timeout: 10 +output.logstash: + hosts: ["%s"]%s +logging.level: warning +`, httpPort, timeoutSec, lsAddr, bulk) + + cfgFile := filepath.Join(t.TempDir(), "rtbeat.yml") + if err := os.WriteFile(cfgFile, []byte(cfg), 0o600); err != nil { + t.Fatalf("write config: %v", err) + } + + cmd := exec.Command(rtbeatBin, "-c", cfgFile, "-e") + buf := &syncBuffer{} + cmd.Stdout = buf + cmd.Stderr = buf + if err := cmd.Start(); err != nil { + t.Fatalf("start rtbeat: %v", err) + } + + rp := &rtbeatProc{cmd: cmd, httpPort: httpPort, stderr: buf, waitErr: make(chan error, 1)} + go func() { rp.waitErr <- rp.cmd.Wait() }() + + if err := waitForHTTP(fmt.Sprintf("http://127.0.0.1:%d/metrics", httpPort), 10*time.Second); err != nil { + rp.kill() + t.Fatalf("rtbeat did not become ready: %v\n--- rtbeat output ---\n%s", err, buf.String()) + } + return rp +} + +func (rp *rtbeatProc) postIn(t *testing.T, body string) (int, string) { + t.Helper() + resp, err := http.Post(fmt.Sprintf("http://127.0.0.1:%d/in", rp.httpPort), + "application/json", strings.NewReader(body)) + if err != nil { + t.Fatalf("POST /in: %v\n--- rtbeat output ---\n%s", err, rp.stderr.String()) + } + defer func() { _ = resp.Body.Close() }() + b, _ := io.ReadAll(resp.Body) + return resp.StatusCode, string(b) +} + +// sigterm sends SIGTERM and returns the channel delivering the process exit +// error (owned by the single Wait goroutine started in startRtbeatCfg). +func (rp *rtbeatProc) sigterm() <-chan error { + _ = rp.cmd.Process.Signal(syscall.SIGTERM) + return rp.waitErr +} + +// kill force-terminates and reaps via the single Wait goroutine. +func (rp *rtbeatProc) kill() { + _ = rp.cmd.Process.Signal(syscall.SIGKILL) + select { + case <-rp.waitErr: + case <-time.After(5 * time.Second): + } +} + +func TestE2EAckAfterDelivery(t *testing.T) { + ls := startLumber(t, ackImmediate, 0) + defer ls.close() + rp := startRtbeat(t, ls.addr, 5) + defer rp.kill() + + code, body := rp.postIn(t, validBatch) + if code != http.StatusOK { + t.Fatalf("status = %d, want 200; body=%s\n--- rtbeat ---\n%s", code, body, rp.stderr.String()) + } + if !waitFor(func() bool { return ls.received() >= 1 }, 5*time.Second) { + t.Fatalf("output did not receive the event (got %d)", ls.received()) + } +} + +// TestE2EWaitsForAckBeforeOK positively proves the 200 is gated on the ack: +// the output delays its ack, and the response must not arrive before then. +func TestE2EWaitsForAckBeforeOK(t *testing.T) { + const delay = 800 * time.Millisecond + ls := startLumber(t, ackDelayed, delay) + defer ls.close() + rp := startRtbeat(t, ls.addr, 5) + defer rp.kill() + + start := time.Now() + code, body := rp.postIn(t, validBatch) + elapsed := time.Since(start) + + if code != http.StatusOK { + t.Fatalf("status = %d, want 200; body=%s\n--- rtbeat ---\n%s", code, body, rp.stderr.String()) + } + if elapsed < delay { + t.Errorf("200 returned after %s, before the %s ack delay — handler acked on receipt, not delivery", elapsed, delay) + } +} + +func TestE2EStallReturns504(t *testing.T) { + ls := startLumber(t, ackNever, 0) + defer ls.close() + rp := startRtbeat(t, ls.addr, 2) // short ack budget to keep the test quick + defer rp.kill() + + start := time.Now() + code, body := rp.postIn(t, validBatch) + elapsed := time.Since(start) + + if code != http.StatusGatewayTimeout { + t.Fatalf("status = %d, want 504; body=%s\n--- rtbeat ---\n%s", code, body, rp.stderr.String()) + } + // It must have actually waited for the (never-arriving) ack, not failed + // fast (a connection error or parse error would return sub-second). + if elapsed < time.Second { + t.Errorf("returned 504 after %s, expected to wait ~ack timeout (2s)", elapsed) + } + // The event did reach the output; it is in-flight and unacked, so the + // sender (which got a 504) retries rather than dropping it. + if !waitFor(func() bool { return ls.received() >= 1 }, 5*time.Second) { + t.Fatalf("output did not receive the in-flight event") + } +} + +// TestE2EMultiMessageFragmentedAck sends a multi-message batch with +// bulk_max_size=1, so the output ships each event as its own lumberjack window +// and acks them separately. The single /in batch must wait for ALL fragment +// acks before returning 200 — exercising the per-batch ack correlation under +// fragmented, multi-callback acks over the wire. +func TestE2EMultiMessageFragmentedAck(t *testing.T) { + const n = 20 + ls := startLumber(t, ackImmediate, 0) + defer ls.close() + rp := startRtbeatCfg(t, ls.addr, 5, 1) + defer rp.kill() + + code, body := rp.postIn(t, makeBatch(n)) + if code != http.StatusOK { + t.Fatalf("status = %d, want 200; body=%s\n--- rtbeat ---\n%s", code, body, rp.stderr.String()) + } + if !waitFor(func() bool { return ls.received() >= n }, 5*time.Second) { + t.Fatalf("output received %d events, want >= %d", ls.received(), n) + } +} + +// TestE2EDrainsInFlightOnShutdown proves a request that is in-flight when +// SIGTERM arrives is still delivered (graceful drain), not dropped. +func TestE2EDrainsInFlightOnShutdown(t *testing.T) { + // The output acks ~1s after receiving, so the request is still in-flight + // once we confirm the event reached the output and then signal shutdown. + ls := startLumber(t, ackDelayed, time.Second) + defer ls.close() + rp := startRtbeat(t, ls.addr, 5) + + codeCh := make(chan int, 1) + go func() { + code, _ := rp.postIn(t, validBatch) + codeCh <- code + }() + + // Gate on the event actually reaching the output (more robust than a bare + // sleep): at this point its ack is still pending, so SIGTERM lands mid-flight. + if !waitFor(func() bool { return ls.received() >= 1 }, 5*time.Second) { + rp.kill() + t.Fatal("event never reached the output before shutdown") + } + exit := rp.sigterm() + + select { + case code := <-codeCh: + if code != http.StatusOK { + t.Fatalf("in-flight request during shutdown returned %d, want 200 (drained)\n--- rtbeat ---\n%s", code, rp.stderr.String()) + } + case <-time.After(8 * time.Second): + rp.kill() + t.Fatal("in-flight request did not complete during graceful shutdown") + } + + select { + case err := <-exit: + if err != nil { + t.Errorf("rtbeat exited with error after SIGTERM: %v\n--- rtbeat ---\n%s", err, rp.stderr.String()) + } + case <-time.After(12 * time.Second): + rp.kill() + t.Fatal("rtbeat did not exit after SIGTERM within the shutdown budget") + } +} + +// --- small helpers --- + +func makeBatch(n int) string { + var sb strings.Builder + fmt.Fprintf(&sb, `{"uuid":"e2e","size":%d,"messages":[`, n) + for i := 0; i < n; i++ { + if i > 0 { + sb.WriteByte(',') + } + fmt.Fprintf(&sb, `{"seq":"%d","producer":"test","payload":{"i":%d}}`, i, i) + } + sb.WriteString(`]}`) + return sb.String() +} + +func freePort(t *testing.T) int { + t.Helper() + l, err := net.Listen("tcp", "127.0.0.1:0") + if err != nil { + t.Fatalf("reserve port: %v", err) + } + defer func() { _ = l.Close() }() + return l.Addr().(*net.TCPAddr).Port +} + +func waitForHTTP(url string, timeout time.Duration) error { + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + resp, err := http.Get(url) + if err == nil { + _ = resp.Body.Close() + return nil + } + time.Sleep(50 * time.Millisecond) + } + return fmt.Errorf("%s not ready within %s", url, timeout) +} + +func waitFor(cond func() bool, timeout time.Duration) bool { + deadline := time.Now().Add(timeout) + for time.Now().Before(deadline) { + if cond() { + return true + } + time.Sleep(20 * time.Millisecond) + } + return cond() +} + +// syncBuffer is a goroutine-safe buffer for capturing subprocess output. +type syncBuffer struct { + mu sync.Mutex + buf []byte +} + +func (b *syncBuffer) Write(p []byte) (int, error) { + b.mu.Lock() + defer b.mu.Unlock() + b.buf = append(b.buf, p...) + return len(p), nil +} + +func (b *syncBuffer) String() string { + b.mu.Lock() + defer b.mu.Unlock() + return string(b.buf) +}