v0.12.2 config validation updated, distributed to their own pacakges and plugins
This commit is contained in:
@@ -63,9 +63,13 @@ func NewFlow(cfg *config.FlowConfig, logger *log.Logger) (*Flow, error) {
|
||||
}
|
||||
f.formatter = formatter
|
||||
|
||||
// Create heartbeat generator with the same formatter
|
||||
if cfg.Heartbeat != nil && cfg.Heartbeat.Enabled {
|
||||
f.heartbeat = NewHeartbeatGenerator(cfg.Heartbeat, formatter, logger)
|
||||
// Create heartbeat generator with the same formatter if configured
|
||||
if cfg.Heartbeat != nil {
|
||||
hb, err := NewHeartbeatGenerator(cfg.Heartbeat, formatter, logger)
|
||||
if err != nil {
|
||||
return nil, fmt.Errorf("heartbeat: %w", err)
|
||||
}
|
||||
f.heartbeat = hb
|
||||
}
|
||||
|
||||
logger.Info("msg", "Flow processor created",
|
||||
|
||||
@@ -3,17 +3,25 @@ package flow
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"logwisp/src/internal/format"
|
||||
"fmt"
|
||||
"sync/atomic"
|
||||
"time"
|
||||
|
||||
"logwisp/src/internal/config"
|
||||
"logwisp/src/internal/core"
|
||||
"logwisp/src/internal/format"
|
||||
|
||||
lconfig "github.com/lixenwraith/config"
|
||||
"github.com/lixenwraith/log"
|
||||
"github.com/lixenwraith/log/formatter"
|
||||
)
|
||||
|
||||
const (
|
||||
MinHeartbeatIntervalMS = 100
|
||||
DefaultHeartbeatIntervalMS = 1000
|
||||
DefaultHeartbeatFormat = "txt"
|
||||
)
|
||||
|
||||
// HeartbeatGenerator produces periodic heartbeat events
|
||||
type HeartbeatGenerator struct {
|
||||
config *config.HeartbeatConfig
|
||||
@@ -24,14 +32,35 @@ type HeartbeatGenerator struct {
|
||||
}
|
||||
|
||||
// NewHeartbeatGenerator creates a new heartbeat generator
|
||||
func NewHeartbeatGenerator(cfg *config.HeartbeatConfig, formatter format.Formatter, logger *log.Logger) *HeartbeatGenerator {
|
||||
func NewHeartbeatGenerator(cfg *config.HeartbeatConfig, formatter format.Formatter, logger *log.Logger) (*HeartbeatGenerator, error) {
|
||||
if cfg == nil || !cfg.Enabled {
|
||||
return nil, nil
|
||||
}
|
||||
|
||||
// Validate
|
||||
if cfg.IntervalMS == 0 {
|
||||
cfg.IntervalMS = DefaultHeartbeatIntervalMS
|
||||
} else if cfg.IntervalMS < MinHeartbeatIntervalMS {
|
||||
return nil, fmt.Errorf("interval_ms: must be >= %d, got %d", MinHeartbeatIntervalMS, cfg.IntervalMS)
|
||||
}
|
||||
|
||||
validateFormat := lconfig.OneOf("txt", "json", "raw", "")
|
||||
if err := validateFormat(cfg.Format); err != nil {
|
||||
return nil, fmt.Errorf("format: %w", err)
|
||||
}
|
||||
|
||||
// Defaults
|
||||
if cfg.Format == "" {
|
||||
cfg.Format = DefaultHeartbeatFormat
|
||||
}
|
||||
|
||||
hg := &HeartbeatGenerator{
|
||||
config: cfg,
|
||||
formatter: formatter,
|
||||
logger: logger,
|
||||
}
|
||||
hg.lastBeat.Store(time.Time{})
|
||||
return hg
|
||||
return hg, nil
|
||||
}
|
||||
|
||||
// Start begins generating heartbeat events
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
package flow
|
||||
|
||||
import (
|
||||
"fmt"
|
||||
"strings"
|
||||
"sync/atomic"
|
||||
|
||||
@@ -8,6 +9,7 @@ import (
|
||||
"logwisp/src/internal/core"
|
||||
"logwisp/src/internal/tokenbucket"
|
||||
|
||||
lconfig "github.com/lixenwraith/config"
|
||||
"github.com/lixenwraith/log"
|
||||
)
|
||||
|
||||
@@ -25,21 +27,36 @@ type RateLimiter struct {
|
||||
|
||||
// NewRateLimiter creates a new pipeline-level rate limiter from configuration
|
||||
func NewRateLimiter(cfg config.RateLimitConfig, logger *log.Logger) (*RateLimiter, error) {
|
||||
// Rate <= 0 means disabled
|
||||
if cfg.Rate <= 0 {
|
||||
return nil, nil // No rate limit
|
||||
}
|
||||
|
||||
// Validate
|
||||
if err := lconfig.NonNegative(cfg.Rate); err != nil {
|
||||
return nil, fmt.Errorf("rate: %w", err)
|
||||
}
|
||||
if err := lconfig.NonNegative(cfg.Burst); err != nil {
|
||||
return nil, fmt.Errorf("burst: %w", err)
|
||||
}
|
||||
if err := lconfig.NonNegative(cfg.MaxEntrySizeBytes); err != nil {
|
||||
return nil, fmt.Errorf("max_entry_size_bytes: %w", err)
|
||||
}
|
||||
|
||||
// Defaults
|
||||
burst := cfg.Burst
|
||||
if burst <= 0 {
|
||||
burst = cfg.Rate // Default burst to rate
|
||||
burst = cfg.Rate
|
||||
}
|
||||
|
||||
var policy config.RateLimitPolicy
|
||||
switch strings.ToLower(cfg.Policy) {
|
||||
case "drop":
|
||||
policy = config.PolicyDrop
|
||||
default:
|
||||
case "pass", "":
|
||||
policy = config.PolicyPass
|
||||
default:
|
||||
return nil, fmt.Errorf("policy: must be one of [drop, pass], got %s", cfg.Policy)
|
||||
}
|
||||
|
||||
l := &RateLimiter{
|
||||
|
||||
Reference in New Issue
Block a user