MM-27744 disable Zap for unit tests. (#15398)
MM-27744 disable Zap for unit tests. Zap has no concept of shutdown or close. Zap is only shutdown when the app exits. Not a problem for console logging, but when creating a new Zap logger that outputs to files on every unit test, that leaves no easy way to clean up until process exit. Depending on what else is running this can exhaust all file handles and cause unit tests to fail. Zap is now disabled unit tests and uses Logr instead, regardless of config settings. `make test-server` peak file handle usage dropped from ~5K to less than 100.
Этот коммит содержится в:
274
vendor/github.com/mattermost/logr/logr.go
сгенерированный
поставляемый
274
vendor/github.com/mattermost/logr/logr.go
сгенерированный
поставляемый
@@ -27,12 +27,13 @@ type Logr struct {
|
||||
shutdown bool
|
||||
lvlCache levelCache
|
||||
|
||||
metricsOnce sync.Once
|
||||
metricsDone chan struct{}
|
||||
metrics MetricsCollector
|
||||
queueSizeGauge Gauge
|
||||
loggedCounter Counter
|
||||
errorCounter Counter
|
||||
metricsInitOnce sync.Once
|
||||
metricsCloseOnce sync.Once
|
||||
metricsDone chan struct{}
|
||||
metrics MetricsCollector
|
||||
queueSizeGauge Gauge
|
||||
loggedCounter Counter
|
||||
errorCounter Counter
|
||||
|
||||
bufferPool sync.Pool
|
||||
|
||||
@@ -114,42 +115,33 @@ func (logr *Logr) Configure(config *cfg.Config) error {
|
||||
return fmt.Errorf("not implemented yet")
|
||||
}
|
||||
|
||||
// AddTarget adds a target to the logger which will receive
|
||||
// log records for outputting.
|
||||
func (logr *Logr) AddTarget(target Target) error {
|
||||
logr.mux.Lock()
|
||||
defer logr.mux.Unlock()
|
||||
|
||||
if logr.shutdown {
|
||||
return fmt.Errorf("logr shut down")
|
||||
}
|
||||
|
||||
logr.tmux.Lock()
|
||||
defer logr.tmux.Unlock()
|
||||
logr.targets = append(logr.targets, target)
|
||||
|
||||
var err error
|
||||
if logr.metrics != nil {
|
||||
if tm, ok := target.(TargetWithMetrics); ok {
|
||||
err = tm.EnableMetrics(logr.metrics, logr.MetricsUpdateFreqMillis)
|
||||
}
|
||||
}
|
||||
|
||||
func (logr *Logr) ensureInit() {
|
||||
logr.once.Do(func() {
|
||||
defer func() {
|
||||
go logr.start()
|
||||
}()
|
||||
|
||||
logr.mux.Lock()
|
||||
defer logr.mux.Unlock()
|
||||
|
||||
logr.maxQueueSizeActual = logr.MaxQueueSize
|
||||
if logr.maxQueueSizeActual == 0 {
|
||||
logr.maxQueueSizeActual = DefaultMaxQueueSize
|
||||
}
|
||||
|
||||
if logr.maxQueueSizeActual < 0 {
|
||||
logr.maxQueueSizeActual = 0
|
||||
}
|
||||
|
||||
logr.in = make(chan *LogRec, logr.maxQueueSizeActual)
|
||||
logr.done = make(chan struct{})
|
||||
|
||||
if logr.UseSyncMapLevelCache {
|
||||
logr.lvlCache = &syncMapLevelCache{}
|
||||
} else {
|
||||
logr.lvlCache = &arrayLevelCache{}
|
||||
}
|
||||
|
||||
if logr.MaxPooledBuffer == 0 {
|
||||
logr.MaxPooledBuffer = DefaultMaxPooledBuffer
|
||||
}
|
||||
@@ -158,11 +150,41 @@ func (logr *Logr) AddTarget(target Target) error {
|
||||
return new(bytes.Buffer)
|
||||
},
|
||||
}
|
||||
|
||||
logr.lvlCache.setup()
|
||||
go logr.start()
|
||||
})
|
||||
logr.resetLevelCache()
|
||||
return err
|
||||
}
|
||||
|
||||
// AddTarget adds one or more targets to the logger which will receive
|
||||
// log records for outputting.
|
||||
func (logr *Logr) AddTarget(targets ...Target) error {
|
||||
if logr.IsShutdown() {
|
||||
return fmt.Errorf("AddTarget called after Logr shut down")
|
||||
}
|
||||
|
||||
logr.ensureInit()
|
||||
metrics := logr.getMetricsCollector()
|
||||
defer logr.ResetLevelCache() // call this after tmux is released
|
||||
|
||||
logr.tmux.Lock()
|
||||
defer logr.tmux.Unlock()
|
||||
|
||||
errs := merror.New()
|
||||
for _, t := range targets {
|
||||
if t == nil {
|
||||
continue
|
||||
}
|
||||
|
||||
logr.targets = append(logr.targets, t)
|
||||
if metrics != nil {
|
||||
if tm, ok := t.(TargetWithMetrics); ok {
|
||||
if err := tm.EnableMetrics(metrics, logr.MetricsUpdateFreqMillis); err != nil {
|
||||
errs.Append(err)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
return errs.ErrorOrNil()
|
||||
}
|
||||
|
||||
// NewLogger creates a Logger using defaults. A `Logger` is light-weight
|
||||
@@ -178,28 +200,13 @@ var levelStatusDisabled = LevelStatus{}
|
||||
// IsLevelEnabled returns true if at least one target has the specified
|
||||
// level enabled. The result is cached so that subsequent checks are fast.
|
||||
func (logr *Logr) IsLevelEnabled(lvl Level) LevelStatus {
|
||||
// Check cache. lvlCache may still be nil if no targets added.
|
||||
if logr.lvlCache == nil {
|
||||
return levelStatusDisabled
|
||||
}
|
||||
status, ok := logr.lvlCache.get(lvl.ID)
|
||||
status, ok := logr.isLevelEnabledFromCache(lvl)
|
||||
if ok {
|
||||
return status
|
||||
}
|
||||
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
|
||||
// Don't accept new log records after shutdown.
|
||||
if logr.shutdown {
|
||||
return levelStatusDisabled
|
||||
}
|
||||
|
||||
status = LevelStatus{}
|
||||
|
||||
// Check each target.
|
||||
logr.tmux.RLock()
|
||||
defer logr.tmux.RUnlock()
|
||||
for _, t := range logr.targets {
|
||||
e, s := t.IsLevelEnabled(lvl)
|
||||
if e {
|
||||
@@ -210,15 +217,45 @@ func (logr *Logr) IsLevelEnabled(lvl Level) LevelStatus {
|
||||
}
|
||||
}
|
||||
}
|
||||
logr.tmux.RUnlock()
|
||||
|
||||
// Cache and return the result.
|
||||
if err := logr.lvlCache.put(lvl.ID, status); err != nil {
|
||||
if err := logr.updateLevelCache(lvl.ID, status); err != nil {
|
||||
logr.ReportError(err)
|
||||
return LevelStatus{}
|
||||
}
|
||||
return status
|
||||
}
|
||||
|
||||
func (logr *Logr) isLevelEnabledFromCache(lvl Level) (LevelStatus, bool) {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
|
||||
// Don't accept new log records after shutdown.
|
||||
if logr.shutdown {
|
||||
return levelStatusDisabled, true
|
||||
}
|
||||
|
||||
// Check cache. lvlCache may still be nil if no targets added.
|
||||
if logr.lvlCache == nil {
|
||||
return levelStatusDisabled, true
|
||||
}
|
||||
status, ok := logr.lvlCache.get(lvl.ID)
|
||||
if ok {
|
||||
return status, true
|
||||
}
|
||||
return LevelStatus{}, false
|
||||
}
|
||||
|
||||
func (logr *Logr) updateLevelCache(id LevelID, status LevelStatus) error {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
if logr.lvlCache != nil {
|
||||
return logr.lvlCache.put(id, status)
|
||||
}
|
||||
return nil
|
||||
}
|
||||
|
||||
// HasTargets returns true only if at least one target exists within the Logr.
|
||||
func (logr *Logr) HasTargets() bool {
|
||||
logr.tmux.RLock()
|
||||
@@ -226,6 +263,71 @@ func (logr *Logr) HasTargets() bool {
|
||||
return len(logr.targets) > 0
|
||||
}
|
||||
|
||||
// TargetInfo provides name and type for a Target.
|
||||
type TargetInfo struct {
|
||||
Name string
|
||||
Type string
|
||||
}
|
||||
|
||||
// TargetInfos enumerates all the targets added to this Logr.
|
||||
// The resulting slice represents a snapshot at time of calling.
|
||||
func (logr *Logr) TargetInfos() []TargetInfo {
|
||||
logr.tmux.RLock()
|
||||
defer logr.tmux.RUnlock()
|
||||
|
||||
infos := make([]TargetInfo, 0)
|
||||
|
||||
for _, t := range logr.targets {
|
||||
inf := TargetInfo{
|
||||
Name: fmt.Sprintf("%v", t),
|
||||
Type: fmt.Sprintf("%T", t),
|
||||
}
|
||||
infos = append(infos, inf)
|
||||
}
|
||||
return infos
|
||||
}
|
||||
|
||||
// RemoveTargets safely removes one or more targets based on the filtering method.
|
||||
// f should return true to delete the target, false to keep it.
|
||||
// When removing a target, best effort is made to write any queued log records before
|
||||
// closing, with cxt determining how much time can be spent in total.
|
||||
// Note, keep the timeout short since this method blocks certain logging operations.
|
||||
func (logr *Logr) RemoveTargets(cxt context.Context, f func(ti TargetInfo) bool) error {
|
||||
var removed bool
|
||||
defer func() {
|
||||
if removed {
|
||||
// call this after tmux is released since
|
||||
// it will lock mux and we don't want to
|
||||
// introduce possible deadlock.
|
||||
logr.ResetLevelCache()
|
||||
}
|
||||
}()
|
||||
|
||||
errs := merror.New()
|
||||
|
||||
logr.tmux.Lock()
|
||||
defer logr.tmux.Unlock()
|
||||
|
||||
cp := make([]Target, 0)
|
||||
|
||||
for _, t := range logr.targets {
|
||||
inf := TargetInfo{
|
||||
Name: fmt.Sprintf("%v", t),
|
||||
Type: fmt.Sprintf("%T", t),
|
||||
}
|
||||
if f(inf) {
|
||||
if err := t.Shutdown(cxt); err != nil {
|
||||
errs.Append(err)
|
||||
}
|
||||
removed = true
|
||||
} else {
|
||||
cp = append(cp, t)
|
||||
}
|
||||
}
|
||||
logr.targets = cp
|
||||
return errs.ErrorOrNil()
|
||||
}
|
||||
|
||||
// ResetLevelCache resets the cached results of `IsLevelEnabled`. This is
|
||||
// called any time a Target is added or a target's level is changed.
|
||||
func (logr *Logr) ResetLevelCache() {
|
||||
@@ -304,15 +406,24 @@ func (logr *Logr) panic(err interface{}) {
|
||||
// timing out. Use `IsTimeoutError` to determine if the returned error is
|
||||
// due to a timeout.
|
||||
func (logr *Logr) Flush() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), logr.flushTimeout())
|
||||
defer cancel()
|
||||
return logr.FlushWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// Flush blocks while flushing the logr queue and all target queues, by
|
||||
// writing existing log records to valid targets.
|
||||
// Any attempts to add new log records will block until flush is complete.
|
||||
// Use `IsTimeoutError` to determine if the returned error is
|
||||
// due to a timeout.
|
||||
func (logr *Logr) FlushWithTimeout(ctx context.Context) error {
|
||||
if !logr.HasTargets() {
|
||||
return nil
|
||||
}
|
||||
|
||||
logr.mux.Lock()
|
||||
defer logr.mux.Unlock()
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), logr.flushTimeout())
|
||||
defer cancel()
|
||||
if logr.IsShutdown() {
|
||||
return errors.New("Flush called on shut down Logr")
|
||||
}
|
||||
|
||||
rec := newFlushLogRec(logr.NewLogger())
|
||||
logr.enqueue(rec)
|
||||
@@ -325,6 +436,15 @@ func (logr *Logr) Flush() error {
|
||||
return nil
|
||||
}
|
||||
|
||||
// IsShutdown returns true if this Logr instance has been shut down.
|
||||
// No further log records can be enqueued and no targets added after
|
||||
// shutdown.
|
||||
func (logr *Logr) IsShutdown() bool {
|
||||
logr.mux.Lock()
|
||||
defer logr.mux.Unlock()
|
||||
return logr.shutdown
|
||||
}
|
||||
|
||||
// Shutdown cleanly stops the logging engine after making best efforts
|
||||
// to flush all targets. Call this function right before application
|
||||
// exit - logr cannot be restarted once shut down.
|
||||
@@ -332,6 +452,17 @@ func (logr *Logr) Flush() error {
|
||||
// timing out. Use `IsTimeoutError` to determine if the returned error is
|
||||
// due to a timeout.
|
||||
func (logr *Logr) Shutdown() error {
|
||||
ctx, cancel := context.WithTimeout(context.Background(), logr.shutdownTimeout())
|
||||
defer cancel()
|
||||
return logr.ShutdownWithTimeout(ctx)
|
||||
}
|
||||
|
||||
// Shutdown cleanly stops the logging engine after making best efforts
|
||||
// to flush all targets. Call this function right before application
|
||||
// exit - logr cannot be restarted once shut down.
|
||||
// Use `IsTimeoutError` to determine if the returned error is due to a
|
||||
// timeout.
|
||||
func (logr *Logr) ShutdownWithTimeout(ctx context.Context) error {
|
||||
logr.mux.Lock()
|
||||
if logr.shutdown {
|
||||
logr.mux.Unlock()
|
||||
@@ -339,16 +470,15 @@ func (logr *Logr) Shutdown() error {
|
||||
}
|
||||
logr.shutdown = true
|
||||
logr.resetLevelCache()
|
||||
if logr.metricsDone != nil {
|
||||
close(logr.metricsDone)
|
||||
logr.metricsDone = nil
|
||||
}
|
||||
logr.mux.Unlock()
|
||||
|
||||
errs := merror.New()
|
||||
logr.metricsCloseOnce.Do(func() {
|
||||
if logr.metricsDone != nil {
|
||||
close(logr.metricsDone)
|
||||
}
|
||||
})
|
||||
|
||||
ctx, cancel := context.WithTimeout(context.Background(), logr.shutdownTimeout())
|
||||
defer cancel()
|
||||
errs := merror.New()
|
||||
|
||||
// close the incoming channel and wait for read loop to exit.
|
||||
if logr.in != nil {
|
||||
@@ -377,9 +507,8 @@ func (logr *Logr) Shutdown() error {
|
||||
// If `OnLoggerError` is not nil, it is called with the error, otherwise the error is
|
||||
// output to `os.Stderr`.
|
||||
func (logr *Logr) ReportError(err interface{}) {
|
||||
if logr.errorCounter != nil {
|
||||
logr.errorCounter.Inc()
|
||||
}
|
||||
logr.incErrorCounter()
|
||||
|
||||
if logr.OnLoggerError == nil {
|
||||
fmt.Fprintln(os.Stderr, err)
|
||||
return
|
||||
@@ -455,7 +584,7 @@ func (logr *Logr) start() {
|
||||
// logr is closed.
|
||||
func (logr *Logr) startMetricsUpdater() {
|
||||
for {
|
||||
updateFreq := logr.MetricsUpdateFreqMillis
|
||||
updateFreq := logr.getMetricsUpdateFreqMillis()
|
||||
if updateFreq == 0 {
|
||||
updateFreq = DefMetricsUpdateFreqMillis
|
||||
}
|
||||
@@ -467,13 +596,17 @@ func (logr *Logr) startMetricsUpdater() {
|
||||
case <-logr.metricsDone:
|
||||
return
|
||||
case <-time.After(time.Duration(updateFreq) * time.Millisecond):
|
||||
if logr.queueSizeGauge != nil {
|
||||
logr.queueSizeGauge.Set(float64(len(logr.in)))
|
||||
}
|
||||
logr.setQueueSizeGauge(float64(len(logr.in)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (logr *Logr) getMetricsUpdateFreqMillis() int64 {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
return logr.MetricsUpdateFreqMillis
|
||||
}
|
||||
|
||||
// fanout pushes a LogRec to all targets.
|
||||
func (logr *Logr) fanout(rec *LogRec) {
|
||||
var target Target
|
||||
@@ -484,6 +617,11 @@ func (logr *Logr) fanout(rec *LogRec) {
|
||||
}()
|
||||
|
||||
var logged bool
|
||||
defer func() {
|
||||
if logged {
|
||||
logr.incLoggedCounter() // call this after tmux is released
|
||||
}
|
||||
}()
|
||||
|
||||
logr.tmux.RLock()
|
||||
defer logr.tmux.RUnlock()
|
||||
@@ -493,10 +631,6 @@ func (logr *Logr) fanout(rec *LogRec) {
|
||||
logged = true
|
||||
}
|
||||
}
|
||||
|
||||
if logged && logr.loggedCounter != nil {
|
||||
logr.loggedCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
// flush drains the queue and notifies when done.
|
||||
|
||||
36
vendor/github.com/mattermost/logr/metrics.go
сгенерированный
поставляемый
36
vendor/github.com/mattermost/logr/metrics.go
сгенерированный
поставляемый
@@ -52,6 +52,12 @@ type TargetWithMetrics interface {
|
||||
EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error
|
||||
}
|
||||
|
||||
func (logr *Logr) getMetricsCollector() MetricsCollector {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
return logr.metrics
|
||||
}
|
||||
|
||||
// SetMetricsCollector enables metrics collection by supplying a MetricsCollector.
|
||||
// The MetricsCollector provides counters and gauges that are updated by log targets.
|
||||
func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error {
|
||||
@@ -59,12 +65,14 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error {
|
||||
return errors.New("collector cannot be nil")
|
||||
}
|
||||
|
||||
logr.mux.Lock()
|
||||
logr.metrics = collector
|
||||
logr.queueSizeGauge, _ = collector.QueueSizeGauge("_logr")
|
||||
logr.loggedCounter, _ = collector.LoggedCounter("_logr")
|
||||
logr.errorCounter, _ = collector.ErrorCounter("_logr")
|
||||
logr.mux.Unlock()
|
||||
|
||||
logr.metricsOnce.Do(func() {
|
||||
logr.metricsInitOnce.Do(func() {
|
||||
logr.metricsDone = make(chan struct{})
|
||||
go logr.startMetricsUpdater()
|
||||
})
|
||||
@@ -75,7 +83,7 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error {
|
||||
defer logr.tmux.RUnlock()
|
||||
for _, target := range logr.targets {
|
||||
if tm, ok := target.(TargetWithMetrics); ok {
|
||||
if err := tm.EnableMetrics(logr.metrics, logr.MetricsUpdateFreqMillis); err != nil {
|
||||
if err := tm.EnableMetrics(collector, logr.MetricsUpdateFreqMillis); err != nil {
|
||||
merr.Append(err)
|
||||
}
|
||||
}
|
||||
@@ -83,3 +91,27 @@ func (logr *Logr) SetMetricsCollector(collector MetricsCollector) error {
|
||||
}
|
||||
return merr.ErrorOrNil()
|
||||
}
|
||||
|
||||
func (logr *Logr) setQueueSizeGauge(val float64) {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
if logr.queueSizeGauge != nil {
|
||||
logr.queueSizeGauge.Set(val)
|
||||
}
|
||||
}
|
||||
|
||||
func (logr *Logr) incLoggedCounter() {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
if logr.loggedCounter != nil {
|
||||
logr.loggedCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
func (logr *Logr) incErrorCounter() {
|
||||
logr.mux.RLock()
|
||||
defer logr.mux.RUnlock()
|
||||
if logr.errorCounter != nil {
|
||||
logr.errorCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
100
vendor/github.com/mattermost/logr/target.go
сгенерированный
поставляемый
100
vendor/github.com/mattermost/logr/target.go
сгенерированный
поставляемый
@@ -4,13 +4,14 @@ import (
|
||||
"context"
|
||||
"fmt"
|
||||
"os"
|
||||
"sync"
|
||||
"time"
|
||||
)
|
||||
|
||||
// Target represents a destination for log records such as file,
|
||||
// database, TCP socket, etc.
|
||||
type Target interface {
|
||||
// SetName provides an option name for the target.
|
||||
// SetName provides an optional name for the target.
|
||||
SetName(name string)
|
||||
|
||||
// IsLevelEnabled returns true if this target should emit
|
||||
@@ -39,7 +40,6 @@ type RecordWriter interface {
|
||||
// in your target type, implement `RecordWriter`, and call `(*Basic).Start`.
|
||||
type Basic struct {
|
||||
target Target
|
||||
name string
|
||||
|
||||
filter Filter
|
||||
formatter Formatter
|
||||
@@ -48,6 +48,10 @@ type Basic struct {
|
||||
done chan struct{}
|
||||
w RecordWriter
|
||||
|
||||
mux sync.RWMutex
|
||||
name string
|
||||
|
||||
metrics bool
|
||||
queueSizeGauge Gauge
|
||||
loggedCounter Counter
|
||||
errorCounter Counter
|
||||
@@ -74,12 +78,14 @@ func (b *Basic) Start(target Target, rw RecordWriter, filter Filter, formatter F
|
||||
b.w = rw
|
||||
go b.start()
|
||||
|
||||
if b.queueSizeGauge != nil {
|
||||
if b.hasMetrics() {
|
||||
go b.startMetricsUpdater()
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) SetName(name string) {
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
b.name = name
|
||||
}
|
||||
|
||||
@@ -117,14 +123,10 @@ func (b *Basic) Log(rec *LogRec) {
|
||||
default:
|
||||
handler := lgr.OnTargetQueueFull
|
||||
if handler != nil && handler(b.target, rec, cap(b.in)) {
|
||||
if b.droppedCounter != nil {
|
||||
b.droppedCounter.Inc()
|
||||
}
|
||||
b.incDroppedCounter()
|
||||
return // drop the record
|
||||
}
|
||||
if b.blockedCounter != nil {
|
||||
b.blockedCounter.Inc()
|
||||
}
|
||||
b.incBlockedCounter()
|
||||
|
||||
select {
|
||||
case <-time.After(lgr.enqueueTimeout()):
|
||||
@@ -136,9 +138,14 @@ func (b *Basic) Log(rec *LogRec) {
|
||||
|
||||
// Metrics enables metrics collection using the provided MetricsCollector.
|
||||
func (b *Basic) EnableMetrics(collector MetricsCollector, updateFreqMillis int64) error {
|
||||
name := fmt.Sprintf("%v", b)
|
||||
|
||||
b.mux.Lock()
|
||||
defer b.mux.Unlock()
|
||||
|
||||
b.metrics = true
|
||||
b.metricsUpdateFreqMillis = updateFreqMillis
|
||||
|
||||
name := fmt.Sprintf("%v", b)
|
||||
var err error
|
||||
|
||||
if b.queueSizeGauge, err = collector.QueueSizeGauge(name); err != nil {
|
||||
@@ -159,8 +166,57 @@ func (b *Basic) EnableMetrics(collector MetricsCollector, updateFreqMillis int64
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b *Basic) hasMetrics() bool {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
return b.metrics
|
||||
}
|
||||
|
||||
func (b *Basic) setQueueSizeGauge(val float64) {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
if b.queueSizeGauge != nil {
|
||||
b.queueSizeGauge.Set(val)
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) incLoggedCounter() {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
if b.loggedCounter != nil {
|
||||
b.loggedCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) incErrorCounter() {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
if b.errorCounter != nil {
|
||||
b.errorCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) incDroppedCounter() {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
if b.droppedCounter != nil {
|
||||
b.droppedCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) incBlockedCounter() {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
if b.blockedCounter != nil {
|
||||
b.blockedCounter.Inc()
|
||||
}
|
||||
}
|
||||
|
||||
// String returns a name for this target. Use `SetName` to specify a name.
|
||||
func (b *Basic) String() string {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
|
||||
if b.name != "" {
|
||||
return b.name
|
||||
}
|
||||
@@ -183,12 +239,10 @@ func (b *Basic) start() {
|
||||
} else {
|
||||
err := b.w.Write(rec)
|
||||
if err != nil {
|
||||
if b.errorCounter != nil {
|
||||
b.errorCounter.Inc()
|
||||
}
|
||||
b.incErrorCounter()
|
||||
rec.Logger().Logr().ReportError(err)
|
||||
} else if b.loggedCounter != nil {
|
||||
b.loggedCounter.Inc()
|
||||
} else {
|
||||
b.incLoggedCounter()
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -199,7 +253,7 @@ func (b *Basic) start() {
|
||||
// target is closed.
|
||||
func (b *Basic) startMetricsUpdater() {
|
||||
for {
|
||||
updateFreq := b.metricsUpdateFreqMillis
|
||||
updateFreq := b.getMetricsUpdateFreqMillis()
|
||||
if updateFreq == 0 {
|
||||
updateFreq = DefMetricsUpdateFreqMillis
|
||||
}
|
||||
@@ -211,13 +265,17 @@ func (b *Basic) startMetricsUpdater() {
|
||||
case <-b.done:
|
||||
return
|
||||
case <-time.After(time.Duration(updateFreq) * time.Millisecond):
|
||||
if b.queueSizeGauge != nil {
|
||||
b.queueSizeGauge.Set(float64(len(b.in)))
|
||||
}
|
||||
b.setQueueSizeGauge(float64(len(b.in)))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
func (b *Basic) getMetricsUpdateFreqMillis() int64 {
|
||||
b.mux.RLock()
|
||||
defer b.mux.RUnlock()
|
||||
return b.metricsUpdateFreqMillis
|
||||
}
|
||||
|
||||
// flush drains the queue and notifies when done.
|
||||
func (b *Basic) flush(done chan<- struct{}) {
|
||||
for {
|
||||
@@ -229,9 +287,7 @@ func (b *Basic) flush(done chan<- struct{}) {
|
||||
if rec.flush == nil {
|
||||
err = b.w.Write(rec)
|
||||
if err != nil {
|
||||
if b.errorCounter != nil {
|
||||
b.errorCounter.Inc()
|
||||
}
|
||||
b.incErrorCounter()
|
||||
rec.Logger().Logr().ReportError(err)
|
||||
}
|
||||
}
|
||||
|
||||
Ссылка в новой задаче
Block a user