rozerolog

package module
v0.0.0-...-5128dcf Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Jul 5, 2026 License: Apache-2.0 Imports: 3 Imported by: 0

README

Zerolog Plugin

The Zerolog plugin provides operators for structured logging observables using the Zerolog logging library.

Installation

go get github.com/samber/ro/plugins/observability/zerolog

Operators

Log

Logs all observable events (next, error, complete) using Zerolog at the specified level.

import (
    "github.com/rs/zerolog"
    "github.com/rs/zerolog/log"
    "github.com/samber/ro"
    rozerolog "github.com/samber/ro/plugins/observability/zerolog"
)

// Create a Zerolog logger
logger := log.With().Str("service", "my-app").Logger()

observable := ro.Pipe1(
    ro.Just(1, 2, 3, 4, 5),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[int]())
defer subscription.Unsubscribe()

// Output:
// {"level":"info","message":"ro.Next: 1"}
// {"level":"info","message":"ro.Next: 2"}
// {"level":"info","message":"ro.Next: 3"}
// {"level":"info","message":"ro.Next: 4"}
// {"level":"info","message":"ro.Next: 5"}
// {"level":"info","message":"ro.Complete"}
LogWithNotification

Logs observable events with structured fields, including the value as a log field.

observable := ro.Pipe1(
    ro.Just("Hello", "World", "Golang"),
    rozerolog.LogWithNotification[string](&logger, zerolog.InfoLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[string]())
defer subscription.Unsubscribe()

// Output:
// {"level":"info","value":"Hello","message":"ro.Next"}
// {"level":"info","value":"World","message":"ro.Next"}
// {"level":"info","value":"Golang","message":"ro.Next"}
// {"level":"info","message":"ro.Complete"}
FatalOnError

Logs fatal errors when the observable emits an error.

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.FatalOnError[int](&logger),
)

subscription := observable.Subscribe(
    ro.NewObserver(
        func(value int) {
            // Handle successful value
        },
        func(err error) {
            // Handle error
        },
        func() {
            // Handle completion
        },
    ),
)
defer subscription.Unsubscribe()

Zerolog Levels

The plugin supports all Zerolog log levels:

  • zerolog.DebugLevel
  • zerolog.InfoLevel
  • zerolog.WarnLevel
  • zerolog.ErrorLevel
  • zerolog.FatalLevel
  • zerolog.PanicLevel
  • zerolog.Disabled
// Log at debug level
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.DebugLevel),
)

Context Support

All operators support context for structured logging:

ctx := context.WithValue(context.Background(), "request_id", "12345")

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)

subscription := observable.SubscribeWithContext(ctx, ro.NoopObserver[int]())
defer subscription.Unsubscribe()

// Output:
// {"level":"info","message":"ro.Next: 1"}
// {"level":"info","message":"ro.Next: 2"}
// {"level":"info","message":"ro.Next: 3"}
// {"level":"info","message":"ro.Complete"}

Structured Data

The plugin works well with structured data types:

type User struct {
    Name string
    Age  int
}

observable := ro.Pipe1(
    ro.Just(
        User{Name: "Alice", Age: 30},
        User{Name: "Bob", Age: 25},
        User{Name: "Charlie", Age: 35},
    ),
    rozerolog.LogWithNotification[User](&logger, zerolog.InfoLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[User]())
defer subscription.Unsubscribe()

// Output:
// {"level":"info","value":"{Alice 30}","message":"ro.Next"}
// {"level":"info","value":"{Bob 25}","message":"ro.Next"}
// {"level":"info","value":"{Charlie 35}","message":"ro.Next"}
// {"level":"info","message":"ro.Complete"}

Logger Configuration

You can configure the Zerolog logger with various options:

Global Logger
// Use the global logger
logger := log.With().Str("service", "my-app").Logger()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)
Console Logger
// Create a console logger
logger := zerolog.New(os.Stdout).With().Timestamp().Logger()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)
JSON Logger
// Create a JSON logger
logger := zerolog.New(os.Stdout).With().Timestamp().Logger()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)
Custom Logger
// Create a custom logger with fields
logger := zerolog.New(os.Stdout).
    With().
    Str("service", "my-app").
    Str("version", "1.0.0").
    Timestamp().
    Logger()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.InfoLevel),
)

Error Handling

The plugin provides different error handling strategies:

Log Errors
// Log errors at a specific level
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.Log[int](&logger, zerolog.ErrorLevel),
)
Fatal on Error
// Log fatal errors and terminate
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.FatalOnError[int](&logger),
)
Structured Error Logging
// Log errors with structured fields
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozerolog.LogWithNotification[int](&logger, zerolog.ErrorLevel),
)

Real-world Example

Here's a practical example that logs API requests:

import (
    "context"
    "github.com/rs/zerolog"
    "github.com/rs/zerolog/log"
    "github.com/samber/ro"
    rozerolog "github.com/samber/ro/plugins/observability/zerolog"
)

// Create a logger for API requests
logger := log.With().
    Str("service", "api-gateway").
    Str("environment", "production").
    Logger()

// Process API requests with logging
pipeline := ro.Pipe2(
    // Source: API requests
    ro.Just("GET /users", "POST /users", "GET /users/123"),
    // Log each request
    rozerolog.LogWithNotification[string](&logger, zerolog.InfoLevel),
)

subscription := pipeline.Subscribe(
    ro.NewObserver(
        func(request string) {
            // Process the request
        },
        func(err error) {
            // Handle error
        },
        func() {
            // Handle completion
        },
    ),
)
defer subscription.Unsubscribe()

Performance Considerations

  • The plugin uses Zerolog's efficient logging mechanisms
  • Context propagation adds minimal overhead
  • Structured logging with fields is optimized
  • Consider log level configuration for production
  • Use appropriate logger configuration for your environment
  • The plugin doesn't block the observable stream
  • Logging is done asynchronously to avoid performance impact
  • Zerolog provides efficient JSON encoding
  • Zero-allocation logging for high-performance applications

Documentation

Index

Examples

Constants

This section is empty.

Variables

This section is empty.

Functions

func FatalOnError

func FatalOnError[T any](logger *zerolog.Logger) func(ro.Observable[T]) ro.Observable[T]

FatalOnError terminates the program with a fatal error when an observable error notification occurs using zerolog logger. Play: https://go.dev/play/p/8na96YvSIci

func Log

func Log[T any](logger *zerolog.Logger, level zerolog.Level) func(ro.Observable[T]) ro.Observable[T]

Log logs all observable notifications (Next, Error, Complete) using zerolog logger with formatted messages. Play: https://go.dev/play/p/wJm9Mw1hkhB

Example
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Log all notifications (Next, Error, Complete)
observable := ro.Pipe1(
	ro.Just(1, 2, 3, 4, 5),
	Log[int](&logger, zerolog.InfoLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[int]())
defer subscription.Unsubscribe()
Output:
{"level":"info","message":"ro.Next: 1"}
{"level":"info","message":"ro.Next: 2"}
{"level":"info","message":"ro.Next: 3"}
{"level":"info","message":"ro.Next: 4"}
{"level":"info","message":"ro.Next: 5"}
{"level":"info","message":"ro.Complete"}
Example (InPipeline)
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Use logging in a complex pipeline
observable := ro.Pipe3(
	ro.Just(1, 2, 3, 4, 5),
	ro.Filter(func(n int) bool { return n%2 == 0 }), // Keep even numbers
	Log[int](&logger, zerolog.InfoLevel),
	ro.Map(func(n int) string { return fmt.Sprintf("Even: %d", n) }),
)

subscription := observable.Subscribe(ro.NoopObserver[string]())
defer subscription.Unsubscribe()
Output:
{"level":"info","message":"ro.Next: 2"}
{"level":"info","message":"ro.Next: 4"}
{"level":"info","message":"ro.Complete"}
Example (WithContext)
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Log with context-aware operations
ctx := context.Background()

observable := ro.Pipe1(
	ro.Just("context", "aware", "logging"),
	LogWithNotification[string](&logger, zerolog.InfoLevel),
)

subscription := observable.SubscribeWithContext(ctx, ro.NoopObserver[string]())
defer subscription.Unsubscribe()
Output:
{"level":"info","value":"context","message":"ro.Next"}
{"level":"info","value":"aware","message":"ro.Next"}
{"level":"info","value":"logging","message":"ro.Next"}
{"level":"info","message":"ro.Complete"}
Example (WithCustomLevels)
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Demonstrate different log levels
observable := ro.Pipe1(
	ro.Just(1, 2, 3, 4, 5),
	Log[int](&logger, zerolog.WarnLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[int]())
defer subscription.Unsubscribe()
Output:
{"level":"warn","message":"ro.Next: 1"}
{"level":"warn","message":"ro.Next: 2"}
{"level":"warn","message":"ro.Next: 3"}
{"level":"warn","message":"ro.Next: 4"}
{"level":"warn","message":"ro.Next: 5"}
{"level":"warn","message":"ro.Complete"}
Example (WithError)
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Log including error notifications
observable := ro.Pipe1(
	ro.NewObservable(func(observer ro.Observer[int]) ro.Teardown {
		observer.Next(1)
		observer.Next(2)
		observer.Error(errors.New("something went wrong"))
		observer.Next(3) // This won't be emitted due to error
		return nil
	}),
	Log[int](&logger, zerolog.ErrorLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[int]())
defer subscription.Unsubscribe()
Output:
{"level":"error","message":"ro.Next: 1"}
{"level":"error","message":"ro.Next: 2"}
{"level":"error","message":"ro.Error: something went wrong"}

func LogWithNotification

func LogWithNotification[T any](logger *zerolog.Logger, level zerolog.Level) func(ro.Observable[T]) ro.Observable[T]

LogWithNotification logs all observable notifications using zerolog logger with structured notification data. Play: https://go.dev/play/p/fHvxhzYzTtV

Example
// Initialize zerolog logger
buff := bufio.NewWriter(os.Stdout)
logger := zerolog.New(buff).With().Logger()
defer buff.Flush()

// Log with structured notification data
observable := ro.Pipe1(
	ro.Just("hello", "world", "golang"),
	LogWithNotification[string](&logger, zerolog.DebugLevel),
)

subscription := observable.Subscribe(ro.NoopObserver[string]())
defer subscription.Unsubscribe()
Output:
{"level":"debug","value":"hello","message":"ro.Next"}
{"level":"debug","value":"world","message":"ro.Next"}
{"level":"debug","value":"golang","message":"ro.Next"}
{"level":"debug","message":"ro.Complete"}

Types

This section is empty.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL