rozap

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: 5 Imported by: 0

README

Zap Plugin

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

Installation

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

Operators

Log

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

import (
    "github.com/samber/ro"
    rozap "github.com/samber/ro/plugins/observability/zap"
    "go.uber.org/zap"
    "go.uber.org/zap/zapcore"
)

// Create a Zap logger
logger, _ := zap.NewProduction()

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

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

logger.Sync()

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

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

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

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

logger.Sync()

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

Logs fatal errors when the observable emits an error.

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.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()

logger.Sync()

Zap Levels

The plugin supports all Zap log levels:

  • zapcore.DebugLevel
  • zapcore.InfoLevel
  • zapcore.WarnLevel
  • zapcore.ErrorLevel
  • zapcore.DPanicLevel
  • zapcore.PanicLevel
  • zapcore.FatalLevel
// Log at debug level
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.Log[int](logger, zapcore.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),
    rozap.Log[int](logger, zapcore.InfoLevel),
)

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

logger.Sync()

// Output:
// {"level":"info","msg":"ro.Next: 1"}
// {"level":"info","msg":"ro.Next: 2"}
// {"level":"info","msg":"ro.Next: 3"}
// {"level":"info","msg":"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},
    ),
    rozap.LogWithNotification[User](logger, zapcore.InfoLevel),
)

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

logger.Sync()

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

Logger Configuration

You can configure the Zap logger with various options:

Production Logger
// Create a production logger
logger, _ := zap.NewProduction()
defer logger.Sync()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.Log[int](logger, zapcore.InfoLevel),
)
Development Logger
// Create a development logger
logger, _ := zap.NewDevelopment()
defer logger.Sync()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.Log[int](logger, zapcore.InfoLevel),
)
Custom Logger
// Create a custom logger
config := zap.NewProductionConfig()
config.OutputPaths = []string{"stdout", "logs/app.log"}
logger, _ := config.Build()
defer logger.Sync()

observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.Log[int](logger, zapcore.InfoLevel),
)
Sugared Logger
// Create a sugared logger
logger, _ := zap.NewProduction()
sugar := logger.Sugar()

// Note: The plugin works with the core logger, not the sugared logger
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.Log[int](logger, zapcore.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),
    rozap.Log[int](logger, zapcore.ErrorLevel),
)
Fatal on Error
// Log fatal errors and terminate
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.FatalOnError[int](logger),
)
Structured Error Logging
// Log errors with structured fields
observable := ro.Pipe1(
    ro.Just(1, 2, 3),
    rozap.LogWithNotification[int](logger, zapcore.ErrorLevel),
)

Real-world Example

Here's a practical example that logs API requests:

import (
    "context"
    "github.com/samber/ro"
    rozap "github.com/samber/ro/plugins/observability/zap"
    "go.uber.org/zap"
    "go.uber.org/zap/zapcore"
)

// Create a logger for API requests
logger, _ := zap.NewProduction()
defer logger.Sync()

// Process API requests with logging
pipeline := ro.Pipe2(
    // Source: API requests
    ro.Just("GET /users", "POST /users", "GET /users/123"),
    // Log each request
    rozap.LogWithNotification[string](logger, zapcore.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 Zap'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
  • Zap provides efficient JSON encoding
  • Remember to call logger.Sync() in production

Documentation

Index

Examples

Constants

This section is empty.

Variables

This section is empty.

Functions

func FatalOnError

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

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

func Log

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

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

Example
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.InfoLevel)

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

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

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	INFO	ro.Next: 1
2024-01-01T12:00:00.000Z	INFO	ro.Next: 2
2024-01-01T12:00:00.000Z	INFO	ro.Next: 3
2024-01-01T12:00:00.000Z	INFO	ro.Next: 4
2024-01-01T12:00:00.000Z	INFO	ro.Next: 5
2024-01-01T12:00:00.000Z	INFO	ro.Complete
Example (InPipeline)
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.DebugLevel)

// 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, zapcore.InfoLevel),
	ro.Map(func(n int) string { return fmt.Sprintf("Even: %d", n) }),
)

subscription := observable.Subscribe(ro.NewObserver(
	func(value string) {
		// Consume values to trigger logging
	},
	func(err error) {
		// Handle errors
	},
	func() {
		// Handle completion
	},
))
defer subscription.Unsubscribe()

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	INFO	ro.Next: 2
2024-01-01T12:00:00.000Z	INFO	ro.Next: 4
2024-01-01T12:00:00.000Z	INFO	ro.Complete
Example (WithContext)
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.DebugLevel)

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

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

subscription := observable.SubscribeWithContext(ctx, ro.NewObserverWithContext(
	func(ctx context.Context, value string) {
		// Consume values to trigger logging
	},
	func(ctx context.Context, err error) {
		// Handle errors
	},
	func(ctx context.Context) {
		// Handle completion
	},
))
defer subscription.Unsubscribe()

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	INFO	ro.Next	{"value": "context"}
2024-01-01T12:00:00.000Z	INFO	ro.Next	{"value": "aware"}
2024-01-01T12:00:00.000Z	INFO	ro.Next	{"value": "logging"}
2024-01-01T12:00:00.000Z	INFO	ro.Complete
Example (WithCustomLevels)
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.DebugLevel)

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

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

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	WARN	ro.Next: 1
2024-01-01T12:00:00.000Z	WARN	ro.Next: 2
2024-01-01T12:00:00.000Z	WARN	ro.Next: 3
2024-01-01T12:00:00.000Z	WARN	ro.Next: 4
2024-01-01T12:00:00.000Z	WARN	ro.Next: 5
2024-01-01T12:00:00.000Z	WARN	ro.Complete
Example (WithError)
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.DebugLevel)

// 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, zapcore.ErrorLevel),
)

subscription := observable.Subscribe(ro.NewObserver(
	func(value int) {
		// Consume values to trigger logging
	},
	func(err error) {
		// Handle errors
	},
	func() {
		// Handle completion
	},
))
defer subscription.Unsubscribe()

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	ERROR	ro.Next: 1
2024-01-01T12:00:00.000Z	ERROR	ro.Next: 2
2024-01-01T12:00:00.000Z	ERROR	ro.Error: something went wrong

func LogWithNotification

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

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

Example
// Initialize zap logger with custom config to match expected output
logger := createTestLogger(zapcore.DebugLevel)

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

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

logger.Sync()
Output:
2024-01-01T12:00:00.000Z	DEBUG	ro.Next	{"value": "hello"}
2024-01-01T12:00:00.000Z	DEBUG	ro.Next	{"value": "world"}
2024-01-01T12:00:00.000Z	DEBUG	ro.Next	{"value": "golang"}
2024-01-01T12:00:00.000Z	DEBUG	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