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
4 changes: 2 additions & 2 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,9 +16,9 @@ Experience SQL-compatible structured log management based on ClickHouse. [Learn
Please let us know at [hello@betterstack.com](mailto:hello@betterstack.com). We're happy to help!

## Credits
`slog-betterstack` was created and maintained by [Samuel Berthe](https://github.com/samber) and released under the MIT license.
`slog-betterstack` was created and maintained by [Samuel Berthe](https://github.com/samber) and released under the MIT license. [Tomáš Procházka](https://github.com/prochac) reported that records logged before exit were lost, and his analysis of the official clients shaped the delivery defaults and the drop accounting.

Thank you, Samuel! ❤️
Thank you, Samuel and Tomáš! ❤️

---

Expand Down
20 changes: 20 additions & 0 deletions doc.go
Original file line number Diff line number Diff line change
@@ -0,0 +1,20 @@
// Package slogbetterstack sends the logs of a Go application to Better Stack through a
// [log/slog] handler.
//
// handler := slogbetterstack.Option{
// Token: "$SOURCE_TOKEN",
// Endpoint: "https://$INGESTING_HOST/",
// Level: slog.LevelInfo, // Debug if omitted
// }.NewBetterstackHandler()
// defer handler.Close()
//
// logger := slog.New(handler)
// logger.Info("Hello from Better Stack!", "service", "UserService")
//
// Records are queued and uploaded in batches by a background goroutine, so logging never
// waits for the network. Close delivers what is still queued and must run before the program
// exits; os.Exit and log.Fatal skip deferred calls. Delivery failures, and a missing token,
// are reported through [Option.OnError], on stderr by default, and counted in [Stats].
//
// See https://betterstack.com/docs/logs/go/ for the full documentation.
package slogbetterstack
20 changes: 14 additions & 6 deletions example-project/main.go
Original file line number Diff line number Diff line change
Expand Up @@ -22,15 +22,19 @@ func main() {
os.Exit(1)
}

option := slogbetterstack.Option{Level: slog.LevelDebug, Token: token}
option := slogbetterstack.Option{
Token: token,
Level: slog.LevelInfo, // Debug if omitted
}
if host := os.Getenv("BETTERSTACK_INGESTING_HOST"); host != "" {
option.Endpoint = "https://" + host + "/"
}

logger := slog.New(option.NewBetterstackHandler())
handler := option.NewBetterstackHandler()
logger := slog.New(handler)
logger = logger.With("release", "v1.0.0")

logger.Debug("Debugging user service.", "service", "UserService")
logger.Info("Starting user service.", "service", "UserService")

logger.With("userID", 123).Error("Unable to fetch user data.")

Expand All @@ -44,8 +48,12 @@ func main() {
With("error", fmt.Errorf("an error")).
Error("a message", slog.Int("count", 1))

// Logs are sent asynchronously: give the handler a moment before the process exits.
time.Sleep(5 * time.Second)
// Records are sent in batches: Close delivers what is still queued before the program exits.
if err := handler.Close(); err != nil {
fmt.Fprintln(os.Stderr, "close:", err)
os.Exit(1)
}

fmt.Println("Sent 3 log records. Open Better Stack → Live tail to see them.")
stats := handler.Stats()
fmt.Printf("Sent %d log records. Open Better Stack → Live tail to see them.\n", stats.Sent)
}
174 changes: 112 additions & 62 deletions handler.go
Original file line number Diff line number Diff line change
@@ -1,10 +1,10 @@
package slogbetterstack

import (
"bytes"
"context"
"encoding/json"
"net/http"
"slices"
"time"

"log/slog"
Expand All @@ -18,40 +18,71 @@ type Option struct {
// log level (default: debug)
Level slog.Leveler

// token
// source token; without it the handler reports the omission through OnError once and
// drops every record instead of sending anything
Token string
// optional: endpoint
Endpoint string
// default: 10s
// optional: how long one upload attempt may take (default: 10s)
Timeout time.Duration

// optional: customize record builder
Converter Converter
// optional: custom marshaler
// optional: custom marshaler, called with the []map[string]any of a batch's records
Marshaler func(v any) ([]byte, error)
// optional: fetch attributes from context
AttrFromContext []func(ctx context.Context) []slog.Attr

// optional: see slog.HandlerOptions
AddSource bool
ReplaceAttr func(groups []string, a slog.Attr) slog.Attr

// Delivery. Records are queued and uploaded in batches by a background goroutine, so
// logging never waits for the network. Call Close before the program exits to deliver
// what is still queued.

// optional: records per upload (default: 1000)
BatchSize int
// optional: how long a partial batch waits before it is uploaded (default: 1s)
BatchInterval time.Duration
// optional: records the queue holds while uploads are behind; further records are
// dropped and counted rather than blocking the application (default: 100000)
MaxQueueSize int
// optional: concurrent uploads (default: 5)
MaxInFlight int
// optional: retries after a failed attempt, for 408, 429, 5xx and network errors;
// negative disables retries (default: 5)
MaxRetries int
// optional: base delay before a retry, doubled on every attempt with jitter; a
// Retry-After header is honoured instead (default: 300ms)
RetryBackoff time.Duration
// optional: how long Close waits for queued and in-flight records (default: 15s)
ShutdownTimeout time.Duration
// optional: send the JSON uncompressed instead of gzip-compressed
DisableCompression bool
// optional: receives every delivery failure and drop summary. It is called from
// background goroutines, possibly several at once, must return promptly and must not
// log through this handler (default: one line on stderr)
OnError func(err error)
// optional: the HTTP client to upload with; Timeout still applies to every request
HTTPClient *http.Client
}

func (o Option) NewBetterstackHandler() slog.Handler {
// NewBetterstackHandler returns a handler that sends records to Better Stack. The handler is
// also a [slog.Handler]; keep the returned value to call Close before the program exits.
// Create one handler per process rather than one per request: each handler owns a goroutine
// and a connection pool from its first record until Close.
func (o Option) NewBetterstackHandler() *BetterstackHandler {
if o.Level == nil {
o.Level = slog.LevelDebug
}

if o.Token == "" {
panic("missing Betterstack token")
}

if o.Endpoint == "" {
o.Endpoint = BetterstackEndpoint
}

if o.Timeout == 0 {
o.Timeout = 10 * time.Second
if o.Timeout <= 0 {
o.Timeout = defaultTimeout
}

if o.Converter == nil {
Expand All @@ -66,43 +97,76 @@ func (o Option) NewBetterstackHandler() slog.Handler {
o.AttrFromContext = []func(ctx context.Context) []slog.Attr{}
}

if o.BatchSize <= 0 {
o.BatchSize = defaultBatchSize
}
if o.BatchInterval <= 0 {
o.BatchInterval = defaultBatchInterval
}
if o.MaxQueueSize <= 0 {
o.MaxQueueSize = defaultMaxQueueSize
}
if o.MaxInFlight <= 0 {
o.MaxInFlight = defaultMaxInFlight
}
switch {
case o.MaxRetries == 0:
o.MaxRetries = defaultMaxRetries
case o.MaxRetries < 0:
o.MaxRetries = 0
}
if o.RetryBackoff <= 0 {
o.RetryBackoff = defaultRetryBackoff
}
if o.ShutdownTimeout <= 0 {
o.ShutdownTimeout = defaultShutdownTimeout
}
if o.OnError == nil {
o.OnError = defaultOnError
}

return &BetterstackHandler{
option: o,
attrs: []slog.Attr{},
groups: []string{},
option: o,
attrs: []slog.Attr{},
groups: []string{},
transport: newTransport(o),
}
}

var _ slog.Handler = (*BetterstackHandler)(nil)

// BetterstackHandler is a [slog.Handler] that sends records to Better Stack. Handlers derived
// with WithAttrs and WithGroup share the queue and the uploads of the handler they came from,
// and Close on any of them closes all of them.
type BetterstackHandler struct {
option Option
attrs []slog.Attr
groups []string
option Option
attrs []slog.Attr
groups []string
transport *transport
}

func (h *BetterstackHandler) Enabled(_ context.Context, level slog.Level) bool {
return level >= h.option.Level.Level()
}

// Handle converts the record and queues it for upload. It never waits for the network: when
// the queue is full the record is dropped and counted. After Close it returns ErrClosed.
func (h *BetterstackHandler) Handle(ctx context.Context, record slog.Record) error {
fromContext := slogcommon.ContextExtractor(ctx, h.option.AttrFromContext)
payload := h.option.Converter(h.option.AddSource, h.option.ReplaceAttr, append(h.attrs, fromContext...), h.groups, &record)

// non-blocking
go func() {
// @TODO: batching ?
_ = send(h.option.Endpoint, h.option.Token, h.option.Timeout, h.option.Marshaler, []map[string]any{payload})
}()
// Every goroutine logging through this handler shares h.attrs, so appending must never
// write into its spare capacity.
attrs := append(slices.Clip(h.attrs), fromContext...)
payload := h.option.Converter(h.option.AddSource, h.option.ReplaceAttr, attrs, h.groups, &record)

return nil
return h.transport.enqueue(payload)
}

func (h *BetterstackHandler) WithAttrs(attrs []slog.Attr) slog.Handler {
return &BetterstackHandler{
option: h.option,
attrs: slogcommon.AppendAttrsToGroup(h.groups, h.attrs, attrs...),
groups: h.groups,
option: h.option,
attrs: slogcommon.AppendAttrsToGroup(h.groups, h.attrs, attrs...),
groups: h.groups,
transport: h.transport,
}
}

Expand All @@ -113,43 +177,29 @@ func (h *BetterstackHandler) WithGroup(name string) slog.Handler {
}

return &BetterstackHandler{
option: h.option,
attrs: h.attrs,
groups: append(h.groups, name),
option: h.option,
attrs: h.attrs,
groups: append(h.groups, name),
transport: h.transport,
}
}

func send(endpoint string, token string, timeout time.Duration, marshaler func(v any) ([]byte, error), payload []map[string]any) error {
client := http.Client{
Timeout: time.Duration(10) * time.Second,
}

json, err := marshaler(payload)
if err != nil {
return err
}

body := bytes.NewBuffer(json)

ctx, cancel := context.WithTimeout(context.Background(), timeout)
defer cancel()

// @TODO: maintain a pool of tcp connections
req, err := http.NewRequestWithContext(ctx, "POST", endpoint, body)
if err != nil {
return err
}

req.Header.Add("authorization", `Bearer `+token)
req.Header.Add("content-type", `application/json`)
req.Header.Add("user-agent", name)

resp, err := client.Do(req)
if err != nil {
return err
}
// Flush uploads every record queued so far and returns once Better Stack has acknowledged
// them, a delivery failed for good, or ctx is done. Failures are reported through OnError.
func (h *BetterstackHandler) Flush(ctx context.Context) error {
return h.transport.flush(ctx)
}

defer resp.Body.Close() //nolint:errcheck
// Close delivers what is still queued, waits for the uploads in flight up to ShutdownTimeout
// and stops the background goroutine. It must run before the program exits: records are
// batched, so without it the last ones are lost. Note that os.Exit and log.Fatal skip deferred
// calls. Close is safe to call more than once; later calls return the first result.
func (h *BetterstackHandler) Close() error {
return h.transport.close()
}

return nil
// Stats reports what happened to the records handed to this handler and the ones derived
// from it.
func (h *BetterstackHandler) Stats() Stats {
return h.transport.stats.snapshot()
}
Loading